use std::io;
#[cfg(feature = "persist")]
use std::path::{Path, PathBuf};
use std::sync::{Arc, RwLock};
#[cfg(feature = "persist")]
use std::time::Instant;
use kevy_hash::KevyHash;
#[cfg(feature = "persist")]
use kevy_persist::reshard::{ShardLayout, commit_reshard, merge_sources, recover_journal};
#[cfg(feature = "persist")]
use kevy_persist::{
Aof, Routing, ShardsMeta, layout, layout::infer_files_n, load_snapshot, read_shards_meta,
replay_aof, write_shards_meta,
};
use kevy_store::Store as Keyspace;
use crate::config::{Config, TtlReaperMode};
#[cfg(feature = "persist")]
use crate::metric::KevyMetric;
use crate::metric::OpenReport;
use crate::store::Inner;
#[inline]
pub(crate) fn shard_idx(key: &[u8], n: usize) -> usize {
if n == 1 {
return 0;
}
let h = key.kevy_hash() as usize;
if n.is_power_of_two() {
h & (n - 1)
} else {
h % n
}
}
#[cfg(feature = "persist")]
struct EmbLayout;
#[cfg(feature = "persist")]
impl ShardLayout for EmbLayout {
fn snapshot_path(&self, dir: &Path, i: usize, _n: usize) -> PathBuf {
layout::snapshot_path(dir, i)
}
fn aof_path(&self, dir: &Path, i: usize, _n: usize) -> PathBuf {
layout::aof_path(dir, i)
}
}
fn fresh_keyspace(config: &Config) -> Keyspace {
let mut s = Keyspace::new();
s.set_max_memory(config.maxmemory, config.eviction_policy);
s.set_cached_clock(matches!(config.ttl_reaper, TtlReaperMode::Background));
s
}
pub(crate) fn build_shards(config: &Config) -> io::Result<(Vec<Arc<RwLock<Inner>>>, OpenReport)> {
let n = config.shards.max(1);
#[allow(unused_mut)] let mut stores: Vec<Keyspace> = (0..n).map(|_| fresh_keyspace(config)).collect();
#[cfg(not(feature = "persist"))]
return Ok((into_inners_mem(stores), OpenReport::default()));
#[cfg(feature = "persist")]
build_shards_persist(config, n, stores)
}
#[cfg(feature = "persist")]
fn build_shards_persist(
config: &Config,
n: usize,
mut stores: Vec<Keyspace>,
) -> io::Result<(Vec<Arc<RwLock<Inner>>>, OpenReport)> {
let Some(dir) = config.data_dir.clone() else {
#[cfg(feature = "tier")]
if config.tier_budget.is_some() {
return Err(io::Error::new(
io::ErrorKind::InvalidInput,
"tiering requires a disk data dir (with_persist); a memory-only store has no cold tier",
));
}
return Ok((into_inners(stores, (0..n).map(|_| None).collect()), OpenReport::default()));
};
std::fs::create_dir_all(&dir)?;
recover_journal(&dir, &EmbLayout)?;
enable_tiering(config, &dir, &mut stores)?;
let mut report = load_or_reshard(&dir, config, n, &mut stores)?;
let aofs = open_live_aofs(config, &dir, n, &mut report)?;
Ok((into_inners(stores, aofs), report))
}
#[cfg(feature = "persist")]
fn open_live_aofs(
config: &Config,
dir: &Path,
n: usize,
report: &mut OpenReport,
) -> io::Result<Vec<Option<Aof>>> {
let aofs: Vec<Option<Aof>> = if config.aof {
(0..n)
.map(|i| {
Aof::open_with_repair(
&layout::aof_path(dir, i),
config.appendfsync,
config.replay_resync,
)
.map(Some)
})
.collect::<io::Result<_>>()?
} else {
(0..n).map(|_| None).collect()
};
for aof in aofs.iter().flatten() {
if let Some(q) = aof.open_quarantine() {
report.quarantine_paths.push(q.to_path_buf());
}
}
Ok(aofs)
}
#[cfg(feature = "persist")]
fn enable_tiering(config: &Config, dir: &Path, stores: &mut [Keyspace]) -> io::Result<()> {
#[cfg(all(feature = "tier", not(target_arch = "wasm32")))]
if config.tier_budget.is_some() {
let per_shard = resolve_tier_budget(config, stores.len())?;
for (i, store) in stores.iter_mut().enumerate() {
store.enable_tiering(&dir.join("tier").join(i.to_string()), per_shard)?;
store.set_tier_max_spill(config.max_spill_value);
}
}
#[cfg(not(all(feature = "tier", not(target_arch = "wasm32"))))]
let _ = (config, dir, stores);
Ok(())
}
#[cfg(all(feature = "persist", feature = "tier", not(target_arch = "wasm32")))]
pub(crate) fn resolve_tier_budget(config: &Config, nshards: usize) -> io::Result<u64> {
let spec = config
.tier_budget
.ok_or_else(|| io::Error::new(io::ErrorKind::InvalidInput, "tiering is not configured"))?;
resolve_tier_spec(spec, nshards)
}
#[cfg(all(feature = "tier", not(target_arch = "wasm32")))]
pub(crate) fn resolve_tier_spec(
spec: crate::config::TierBudgetSpec,
nshards: usize,
) -> io::Result<u64> {
use crate::config::TierBudgetSpec;
if let TierBudgetSpec::Percent(p) = spec
&& !(1..=100).contains(&p)
{
return Err(io::Error::new(
io::ErrorKind::InvalidInput,
format!("tiering budget percent must be 1..=100, got {p}"),
));
}
let total = spec.resolve().ok_or_else(|| {
io::Error::new(
io::ErrorKind::Unsupported,
"tiering budget auto/percent: no memory bound detected on this host — \
use with_tier_budget(bytes)",
)
})?;
Ok((total / nshards.max(1) as u64).max(1))
}
#[cfg(all(feature = "tier", not(target_arch = "wasm32")))]
pub(crate) fn tier_tick_upkeep(
g: &mut crate::store::Inner,
spec: Option<crate::config::TierBudgetSpec>,
nshards: usize,
) {
use crate::config::TierBudgetSpec;
if !g.store.tier_enabled() {
return;
}
if let Some(spec @ (TierBudgetSpec::Auto | TierBudgetSpec::Percent(_))) = spec
&& let Ok(per_shard) = resolve_tier_spec(spec, nshards)
{
g.store.set_tier_budget(per_shard);
}
#[cfg(feature = "index")]
let reserved = g.idx_segs.reserved_bytes() + g.view_segs.reserved_bytes();
#[cfg(not(feature = "index"))]
let reserved = 0u64;
g.store.set_tier_reserved(reserved);
}
#[cfg(feature = "persist")]
fn load_or_reshard(
dir: &Path,
config: &Config,
n: usize,
stores: &mut [Keyspace],
) -> io::Result<OpenReport> {
let meta_path = layout::shards_meta_path(dir);
let prev = read_shards_meta(&meta_path);
let same_layout = match prev {
Some(m) => m.n == n && m.routing == Routing::KevyHash,
None => n == 1 && infer_files_n(dir) <= 1,
};
if same_layout {
let report = load_in_place(dir, config, n, stores)?;
write_shards_meta(&meta_path, ShardsMeta { n, routing: Routing::KevyHash })?;
return Ok(report);
}
{
let src_n = prev.map(|m| m.n).or_else(|| {
let k = infer_files_n(dir);
(k > 1).then_some(k)
});
reshard(dir, config, n, src_n, stores)
}
}
#[cfg(feature = "persist")]
fn load_in_place(
dir: &Path,
config: &Config,
_n: usize,
stores: &mut [Keyspace],
) -> io::Result<OpenReport> {
let mut report = OpenReport::default();
let start = Instant::now();
for (i, store) in stores.iter_mut().enumerate() {
let snap = layout::snapshot_path(dir, i);
if snap.exists() {
load_snapshot(store, &snap)?;
}
let aof = layout::aof_path(dir, i);
if aof.exists() {
let mut frames: u64 = 0;
let apply = |args: kevy_persist::Argv| {
crate::replay::apply(store, &args);
frames += 1;
if frames.is_multiple_of(kevy_persist::REPLAY_DEMOTE_INTERVAL) {
store.demote_to_watermark();
}
};
let r = if config.replay_resync {
kevy_persist::replay_aof_resync(&aof, apply)?
} else {
replay_aof(&aof, apply)?
};
report.replayed_commands += r.commands;
report.replayed_bytes += r.replayed_bytes;
report.dropped_bytes += r.dropped_bytes;
report.corrupt |= r.corrupt;
report.resynced_bytes += r.resynced_ranges.iter().map(|(a, b)| b - a).sum::<u64>();
}
store.demote_to_watermark();
}
report.elapsed_ms = start.elapsed().as_millis() as u64;
emit_replay(config, &report);
Ok(report)
}
#[cfg(feature = "persist")]
fn reshard(
dir: &Path,
config: &Config,
n: usize,
prev_n: Option<usize>,
stores: &mut [Keyspace],
) -> io::Result<OpenReport> {
let lay = EmbLayout;
let src_n = prev_n.unwrap_or(1);
let (temp, report) = merge_into_temp(dir, config, src_n)?;
redistribute(&temp, n, stores);
commit_reshard(dir, src_n, ShardsMeta { n, routing: Routing::KevyHash }, stores, &lay)?;
#[cfg(all(feature = "tier", not(target_arch = "wasm32")))]
if config.tier_budget.is_some() {
drop(temp);
let _ = std::fs::remove_dir_all(dir.join("tier").join(".reshard-merge"));
}
Ok(report)
}
#[cfg(feature = "persist")]
fn merge_into_temp(dir: &Path, config: &Config, src_n: usize) -> io::Result<(Keyspace, OpenReport)> {
let lay = EmbLayout;
let mut temp = fresh_keyspace(config);
#[cfg(all(feature = "tier", not(target_arch = "wasm32")))]
if config.tier_budget.is_some() {
let budget = resolve_tier_budget(config, 1)?;
temp.enable_tiering(&dir.join("tier").join(".reshard-merge"), budget)?;
temp.set_tier_max_spill(config.max_spill_value);
}
let mut total_cmds = 0u64;
let start = Instant::now();
merge_sources(dir, src_n, &lay, &mut temp, |store, args| {
total_cmds += 1;
crate::replay::apply(store, &args);
if total_cmds.is_multiple_of(kevy_persist::REPLAY_DEMOTE_INTERVAL) {
store.demote_to_watermark();
}
})?;
let total_bytes = (0..src_n)
.map(|i| lay.aof_path(dir, i, src_n))
.filter_map(|p| std::fs::metadata(p).ok())
.map(|m| m.len())
.sum();
let report = OpenReport {
replayed_commands: total_cmds,
replayed_bytes: total_bytes,
elapsed_ms: start.elapsed().as_millis() as u64,
..OpenReport::default()
};
emit_replay(config, &report);
Ok((temp, report))
}
#[cfg(feature = "persist")]
fn redistribute(temp: &Keyspace, n: usize, stores: &mut [Keyspace]) {
temp.snapshot_each(|key, value, ttl_ms| {
let hot;
let value = match temp.materialize_cold(value) {
Some(v) => {
hot = v;
&hot
}
None => value,
};
let target = &mut stores[shard_idx(key, n)];
target.load_value(key, value, ttl_ms);
target.try_demote_after_write();
});
}
#[cfg(feature = "persist")]
fn emit_replay(config: &Config, report: &OpenReport) {
if let Some(sink) = &config.metric_sink {
sink.emit(KevyMetric::Replay {
commands: report.replayed_commands,
bytes: report.replayed_bytes + report.dropped_bytes,
elapsed_ms: report.elapsed_ms,
dropped_bytes: report.dropped_bytes,
corrupt: report.corrupt,
});
}
}
#[cfg(feature = "persist")]
fn into_inners(stores: Vec<Keyspace>, aofs: Vec<Option<Aof>>) -> Vec<Arc<RwLock<Inner>>> {
stores
.into_iter()
.zip(aofs)
.map(|(store, aof)| Arc::new(RwLock::new(Inner::new(store, aof))))
.collect()
}
#[cfg(not(feature = "persist"))]
fn into_inners_mem(stores: Vec<Keyspace>) -> Vec<Arc<RwLock<Inner>>> {
stores
.into_iter()
.map(|store| Arc::new(RwLock::new(Inner::new(store))))
.collect()
}