use std::path::PathBuf;
use std::time::Duration;
use myelon::MyelonWaitStrategy;
#[derive(Debug, Clone)]
pub struct DaemonConfig {
pub shm_prefixes: Vec<String>,
pub tcp_addrs: Vec<String>,
pub http_addrs: Vec<String>,
pub shm_depth: usize,
pub wait_strategy: MyelonWaitStrategy,
pub tcp_tpc_threads: usize,
pub tcp_dispatch_workers: usize,
pub http_tpc_threads: usize,
pub http_dispatch_workers: usize,
pub puffer_dir: PathBuf,
pub puffer_ram_bytes: Option<u64>,
pub puffer_disk_bytes: Option<u64>,
pub puffer_block_size_bytes: Option<u64>,
pub s3_prefix: String,
pub slatedb_path: PathBuf,
pub namespace: String,
pub namespace_max_bytes: Option<u64>,
pub eviction_interval: Duration,
pub quiet_banner: bool,
pub trace_daemon: bool,
pub prefetch_dry_run: bool,
}
const DAEMON_FALLBACK_PUFFER_DIR: &str = "/tmp/wombatkv-puffer-shm-foyer";
fn default_daemon_puffer_dir() -> PathBuf {
let home = std::env::var("HOME").ok().filter(|s| !s.is_empty());
let candidate = match home {
Some(h) => PathBuf::from(h).join(".wombatkv").join("puffer"),
None => return PathBuf::from(DAEMON_FALLBACK_PUFFER_DIR),
};
match std::fs::create_dir_all(&candidate) {
Ok(()) => candidate,
Err(err) => {
eprintln!(
"wombatkv-daemon: could not create default puffer dir {} \
({err}); falling back to {DAEMON_FALLBACK_PUFFER_DIR}",
candidate.display()
);
PathBuf::from(DAEMON_FALLBACK_PUFFER_DIR)
}
}
}
impl Default for DaemonConfig {
fn default() -> Self {
let puffer_dir = default_daemon_puffer_dir();
let slatedb_path = puffer_dir.join("slatedb");
Self {
shm_prefixes: Vec::new(),
tcp_addrs: Vec::new(),
http_addrs: Vec::new(),
shm_depth: crate::DEFAULT_RING_DEPTH,
wait_strategy: MyelonWaitStrategy::Block,
tcp_tpc_threads: 2,
tcp_dispatch_workers: 8,
http_tpc_threads: 2,
http_dispatch_workers: 8,
puffer_dir,
puffer_ram_bytes: None,
puffer_disk_bytes: None,
puffer_block_size_bytes: None,
s3_prefix: "kv/puffer-shm".to_string(),
slatedb_path,
namespace: "default".to_string(),
namespace_max_bytes: None,
eviction_interval: Duration::from_secs(30),
quiet_banner: false,
trace_daemon: false,
prefetch_dry_run: false,
}
}
}
impl DaemonConfig {
#[must_use]
pub fn from_env() -> Self {
let mut cfg = Self::default();
if let Ok(v) = std::env::var("WMBT_KV_DAEMON_SHM_PREFIX") {
if !v.is_empty() {
cfg.shm_prefixes.push(v);
}
}
if let Ok(v) = std::env::var("WMBT_KV_TCP") {
for chunk in v.split(',') {
let chunk = chunk.trim();
if !chunk.is_empty() {
cfg.tcp_addrs.push(chunk.to_string());
}
}
}
if let Ok(v) = std::env::var("WMBT_KV_HTTP") {
for chunk in v.split(',') {
let chunk = chunk.trim();
if !chunk.is_empty() {
cfg.http_addrs.push(chunk.to_string());
}
}
}
if let Some(n) = env_parse::<usize>("WMBT_KV_DAEMON_SHM_DEPTH") {
cfg.shm_depth = n;
}
cfg.wait_strategy = match std::env::var("WMBT_KV_DAEMON_SHM_WAIT_STRATEGY")
.ok()
.as_deref()
.map(str::trim)
{
Some("busyspin" | "BusySpin" | "spin") => MyelonWaitStrategy::BusySpin,
Some("block" | "Block" | "park") => MyelonWaitStrategy::Block,
Some(other) if !other.is_empty() => {
eprintln!(
"WombatKV: unknown WMBT_KV_DAEMON_SHM_WAIT_STRATEGY={other:?}, defaulting to 'block'"
);
MyelonWaitStrategy::Block
}
_ => MyelonWaitStrategy::Block,
};
if let Some(n) = env_parse::<usize>("WMBT_KV_TCP_TPC_THREADS") {
cfg.tcp_tpc_threads = n.max(1);
}
if let Some(n) = env_parse::<usize>("WMBT_KV_TCP_DISPATCH_WORKERS") {
cfg.tcp_dispatch_workers = n.max(1);
}
if let Some(n) = env_parse::<usize>("WMBT_KV_HTTP_TPC_THREADS") {
cfg.http_tpc_threads = n.max(1);
}
if let Some(n) = env_parse::<usize>("WMBT_KV_HTTP_DISPATCH_WORKERS") {
cfg.http_dispatch_workers = n.max(1);
}
if let Ok(v) = std::env::var("WMBT_KV_PUFFER_DIR") {
if !v.is_empty() {
cfg.puffer_dir = PathBuf::from(&v);
cfg.slatedb_path = cfg.puffer_dir.join("slatedb");
}
}
cfg.puffer_ram_bytes = env_parse::<u64>("WMBT_KV_PUFFER_RAM_BYTES");
cfg.puffer_disk_bytes = env_parse::<u64>("WMBT_KV_PUFFER_DISK_BYTES");
cfg.puffer_block_size_bytes = env_parse::<u64>("WMBT_KV_PUFFER_BLOCK_SIZE_BYTES");
if let Ok(v) = std::env::var("WMBT_KV_S3_PREFIX") {
if !v.is_empty() {
cfg.s3_prefix = v;
}
}
if let Ok(v) = std::env::var("WMBT_KV_SLATEDB_PATH") {
if !v.is_empty() {
cfg.slatedb_path = PathBuf::from(v);
}
}
if let Ok(v) = std::env::var("WMBT_KV_NAMESPACE") {
if !v.is_empty() {
cfg.namespace = v;
}
}
cfg.namespace_max_bytes = env_parse::<u64>("WMBT_KV_NAMESPACE_MAX_BYTES");
if let Some(secs) = env_parse::<u64>("WMBT_KV_EVICTION_INTERVAL_SECS") {
cfg.eviction_interval = Duration::from_secs(secs);
}
cfg.quiet_banner = env_truthy("WMBT_KV_QUIET_BANNER");
cfg.trace_daemon = std::env::var("WMBT_KV_DAEMON_SHM_TRACE_DAEMON").is_ok();
cfg.prefetch_dry_run = env_truthy("WMBT_KV_PREFETCH_DRY_RUN");
cfg
}
}
fn env_parse<T: std::str::FromStr>(key: &str) -> Option<T> {
std::env::var(key).ok().and_then(|v| v.parse::<T>().ok())
}
fn env_truthy(key: &str) -> bool {
matches!(
std::env::var(key).ok().as_deref(),
Some("1" | "true" | "True" | "TRUE" | "yes" | "Yes" | "YES" | "on" | "On" | "ON")
)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn defaults_are_sensible() {
let cfg = DaemonConfig::default();
assert_eq!(cfg.shm_depth, crate::DEFAULT_RING_DEPTH);
assert_eq!(cfg.tcp_tpc_threads, 2);
assert_eq!(cfg.http_tpc_threads, 2);
assert_eq!(cfg.tcp_dispatch_workers, 8);
assert_eq!(cfg.http_dispatch_workers, 8);
assert!(matches!(cfg.wait_strategy, MyelonWaitStrategy::Block));
assert!(!cfg.quiet_banner);
assert!(!cfg.trace_daemon);
assert!(!cfg.prefetch_dry_run);
assert_eq!(cfg.namespace, "default");
assert_eq!(cfg.s3_prefix, "kv/puffer-shm");
assert_eq!(cfg.eviction_interval, Duration::from_secs(30));
}
#[test]
fn env_parse_returns_none_for_garbage() {
let key = "WMBT_KV_TEST_GARBAGE_THAT_NEVER_EXISTS_12345";
assert!(env_parse::<u64>(key).is_none());
}
#[test]
fn env_truthy_matches_canonical_set() {
for tv in ["1", "true", "True", "TRUE", "yes", "on", "ON"] {
assert!(matches!(
Some(tv),
Some("1" | "true" | "True" | "TRUE" | "yes" | "Yes" | "YES" | "on" | "On" | "ON")
));
}
for fv in ["0", "false", "no", "off", ""] {
assert!(!matches!(
Some(fv),
Some("1" | "true" | "True" | "TRUE" | "yes" | "Yes" | "YES" | "on" | "On" | "ON")
));
}
}
}