use std::io;
use std::sync::Arc;
use kevy_persist::{load_snapshot, replay_aof};
use kevy_uring::IoUring;
use crate::Commands;
use crate::shard::Shard;
impl<C: Commands> Shard<C> {
pub(crate) fn prepare_uring_shard(&mut self) -> io::Result<()> {
self.commands.on_shard_start(self.id);
if let Some(rx) = &self.replica_inbox {
rx.attach_waker(Arc::clone(&self.waker));
}
let snap = self.snapshot_path();
if snap.exists()
&& let Err(e) = load_snapshot(&mut self.store, &snap)
{
eprintln!("kevy: shard {} failed to load {}: {e}", self.id, snap.display());
}
if self.aof.is_some() {
let aof_path = self.aof_path();
let commands = &self.commands;
let store = &mut self.store;
let mut frames: u64 = 0;
let apply = |args: kevy_persist::Argv| {
crate::shard_run::replay_dispatch(commands, store, &args);
frames += 1;
if frames.is_multiple_of(kevy_persist::REPLAY_DEMOTE_INTERVAL) {
store.demote_to_watermark();
}
};
let report = if self.replay_resync {
kevy_persist::replay_aof_resync(&aof_path, apply)?
} else {
replay_aof(&aof_path, apply)?
};
self.commands.on_replay_report(report.dropped_bytes, report.corrupt);
}
self.store.demote_to_watermark();
Ok(())
}
}
pub(crate) const URING_ENTRIES: u32 = 2048;
pub(crate) const PBUF_ENTRIES: u16 = 4096;
pub(crate) const PBUF_SIZE: u32 = 16 * 1024;
pub(crate) const PBUF_GROUP: u16 = 0;
pub(crate) fn io_uring_available() -> bool {
match IoUring::new(URING_ENTRIES) {
Ok(ring) => ring
.register_buf_ring(PBUF_ENTRIES, PBUF_SIZE, PBUF_GROUP)
.is_ok(),
Err(_) => false,
}
}
pub(crate) fn build_uring() -> io::Result<(IoUring, kevy_uring::ProvidedBufRing)> {
let sqpoll = matches!(
std::env::var("KEVY_SQPOLL").ok().as_deref(),
Some(v) if !v.is_empty() && v != "0" && v != "off" && v != "no" && v != "false"
);
let ring = if sqpoll {
IoUring::new_sqpoll(URING_ENTRIES, 1000, None)?
} else {
IoUring::new(URING_ENTRIES)?
};
let pbuf = ring.register_buf_ring(PBUF_ENTRIES, PBUF_SIZE, PBUF_GROUP)?;
Ok((ring, pbuf))
}