use std::io::{Read, Write};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use kevy_testnet::free_port;
fn req(parts: &[&[u8]]) -> Vec<u8> {
let mut v = format!("*{}\r\n", parts.len()).into_bytes();
for p in parts {
v.extend_from_slice(format!("${}\r\n", p.len()).as_bytes());
v.extend_from_slice(p);
v.extend_from_slice(b"\r\n");
}
v
}
fn read_reply(s: &mut std::net::TcpStream, expected: &[u8]) {
let mut buf = vec![0u8; expected.len()];
s.read_exact(&mut buf).unwrap();
assert_eq!(
&buf,
expected,
"expected {:?}",
String::from_utf8_lossy(expected)
);
}
fn wait_for(what: &str, mut cond: impl FnMut() -> bool) {
for _ in 0..1000 {
if cond() {
return;
}
std::thread::sleep(std::time::Duration::from_millis(10));
}
panic!("timed out waiting for {what}");
}
fn with_runtime(port: u16, dir: &std::path::Path, nshards: usize, body: impl FnOnce(u16)) {
with_runtime_configured(port, dir, nshards, |rt| rt, body);
}
fn with_runtime_configured<F>(
port: u16,
dir: &std::path::Path,
nshards: usize,
configure: F,
body: impl FnOnce(u16),
) where
F: FnOnce(kevy_rt::Runtime<kevy::KevyCommands>) -> kevy_rt::Runtime<kevy::KevyCommands>
+ Send
+ 'static,
{
let stop = Arc::new(AtomicBool::new(false));
let stop_t = stop.clone();
let dir = dir.to_path_buf();
let handle = std::thread::spawn(move || {
let rt = kevy_rt::Runtime::builder(kevy::KevyCommands::sharded(nshards)).bind([127, 0, 0, 1], port).shards(nshards)
.with_data_dir(dir);
let rt = configure(rt);
rt.run(stop_t).unwrap();
});
let mut up = false;
for _ in 0..200 {
if std::net::TcpStream::connect(("127.0.0.1", port)).is_ok() {
up = true;
break;
}
std::thread::sleep(std::time::Duration::from_millis(5));
}
assert!(up, "runtime did not start");
body(port);
stop.store(true, Ordering::Relaxed);
let _ = handle.join();
}
#[test]
fn data_survives_restart_via_save() {
let dir = std::env::temp_dir().join(format!(
"kevy-persist-{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
std::fs::create_dir_all(&dir).unwrap();
let nshards = 4;
let port = free_port();
with_runtime(port, &dir, nshards, |p| {
let mut c = std::net::TcpStream::connect(("127.0.0.1", p)).unwrap();
for i in 0..100u32 {
c.write_all(&req(&[
b"SET",
format!("k{i}").as_bytes(),
format!("v{i}").as_bytes(),
]))
.unwrap();
read_reply(&mut c, b"+OK\r\n");
}
c.write_all(&req(&[b"SAVE"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
});
let dumps = (0..nshards)
.filter(|i| dir.join(format!("dump-{i}.rdb")).exists())
.count();
assert!(dumps > 0, "no snapshot files were written");
let port2 = free_port();
with_runtime(port2, &dir, nshards, |p| {
let mut c = std::net::TcpStream::connect(("127.0.0.1", p)).unwrap();
for i in 0..100u32 {
c.write_all(&req(&[b"GET", format!("k{i}").as_bytes()]))
.unwrap();
let want = format!("v{i}");
read_reply(
&mut c,
format!("${}\r\n{}\r\n", want.len(), want).as_bytes(),
);
}
});
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn bgrewriteaof_shrinks_log_and_preserves_data() {
let dir = std::env::temp_dir().join(format!(
"kevy-bgrewrite-{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
std::fs::create_dir_all(&dir).unwrap();
let nshards = 4;
let port = free_port();
let mut post_size: u64 = 0;
with_runtime(port, &dir, nshards, |p| {
let mut c = std::net::TcpStream::connect(("127.0.0.1", p)).unwrap();
for i in 0..40u32 {
for rev in 0..50u32 {
c.write_all(&req(&[
b"SET",
format!("k{i}").as_bytes(),
format!("v{i}-r{rev}").as_bytes(),
]))
.unwrap();
read_reply(&mut c, b"+OK\r\n");
}
}
c.write_all(&req(&[b"BGREWRITEAOF"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
let sum_aof = || -> u64 {
(0..nshards)
.map(|s| {
std::fs::metadata(dir.join(format!("aof-{s}.aof")))
.map_or(0, |m| m.len())
})
.sum()
};
wait_for("rewritten AOF to swap in", || sum_aof() < 10_000);
post_size = sum_aof();
assert!(post_size > 0, "rewritten AOF should not be empty");
});
let port2 = free_port();
with_runtime(port2, &dir, nshards, |p| {
let mut c = std::net::TcpStream::connect(("127.0.0.1", p)).unwrap();
for i in 0..40u32 {
c.write_all(&req(&[b"GET", format!("k{i}").as_bytes()]))
.unwrap();
let want = format!("v{i}-r49");
read_reply(
&mut c,
format!("${}\r\n{}\r\n", want.len(), want).as_bytes(),
);
}
});
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn aof_truncated_tail_is_tolerated_on_restart() {
let dir = std::env::temp_dir().join(format!(
"kevy-truncated-{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
std::fs::create_dir_all(&dir).unwrap();
let nshards = 1; let port = free_port();
with_runtime(port, &dir, nshards, |p| {
let mut c = std::net::TcpStream::connect(("127.0.0.1", p)).unwrap();
for i in 0..20u32 {
c.write_all(&req(&[
b"SET",
format!("survivor{i}").as_bytes(),
b"v".to_vec().as_slice(),
]))
.unwrap();
read_reply(&mut c, b"+OK\r\n");
}
c.write_all(&req(&[b"BGREWRITEAOF"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
});
let aof_path = dir.join("aof-0.aof");
let mut bytes = std::fs::read(&aof_path).unwrap();
let prefix_len = bytes.len();
bytes.extend_from_slice(b"*3\r\n$3\r\nSET\r\n$5\r\nfoo");
std::fs::write(&aof_path, &bytes).unwrap();
let corrupted_len = bytes.len();
assert!(corrupted_len > prefix_len, "test should have appended garbage");
let port2 = free_port();
with_runtime(port2, &dir, nshards, |p| {
let mut c = std::net::TcpStream::connect(("127.0.0.1", p)).unwrap();
for i in 0..20u32 {
c.write_all(&req(&[b"GET", format!("survivor{i}").as_bytes()]))
.unwrap();
read_reply(&mut c, b"$1\r\nv\r\n");
}
c.write_all(&req(&[b"GET", b"foo"])).unwrap();
read_reply(&mut c, b"$-1\r\n");
});
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn data_survives_restart_via_aof_without_save() {
let dir = std::env::temp_dir().join(format!(
"kevy-aof-{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
std::fs::create_dir_all(&dir).unwrap();
let nshards = 4;
let port = free_port();
with_runtime(port, &dir, nshards, |p| {
let mut c = std::net::TcpStream::connect(("127.0.0.1", p)).unwrap();
for i in 0..100u32 {
c.write_all(&req(&[
b"SET",
format!("a{i}").as_bytes(),
format!("b{i}").as_bytes(),
]))
.unwrap();
read_reply(&mut c, b"+OK\r\n");
}
for i in 1..=5u32 {
c.write_all(&req(&[b"INCR", b"counter"])).unwrap();
let want = format!(":{i}\r\n");
read_reply(&mut c, want.as_bytes());
}
});
assert!(!dir.join("dump-0.rdb").exists());
let port2 = free_port();
with_runtime(port2, &dir, nshards, |p| {
let mut c = std::net::TcpStream::connect(("127.0.0.1", p)).unwrap();
for i in 0..100u32 {
c.write_all(&req(&[b"GET", format!("a{i}").as_bytes()]))
.unwrap();
let want = format!("b{i}");
read_reply(
&mut c,
format!("${}\r\n{}\r\n", want.len(), want).as_bytes(),
);
}
c.write_all(&req(&[b"GET", b"counter"])).unwrap();
read_reply(&mut c, b"$1\r\n5\r\n");
});
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn restart_tolerates_corrupt_snapshot() {
let dir = std::env::temp_dir().join(format!(
"kevy-corrupt-snap-{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
std::fs::create_dir_all(&dir).unwrap();
std::fs::write(dir.join("dump-0.rdb"), b"NOT A REAL KEVY SNAPSHOT").unwrap();
let port = free_port();
with_runtime(port, &dir, 1, |p| {
let mut c = std::net::TcpStream::connect(("127.0.0.1", p)).unwrap();
c.write_all(&req(&[b"PING"])).unwrap();
read_reply(&mut c, b"+PONG\r\n");
c.write_all(&req(&[b"SET", b"after-corrupt", b"ok"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"GET", b"after-corrupt"])).unwrap();
read_reply(&mut c, b"$2\r\nok\r\n");
});
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn auto_aof_rewrite_fires_when_threshold_crossed() {
let dir = std::env::temp_dir().join(format!(
"kevy-auto-rewrite-{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
std::fs::create_dir_all(&dir).unwrap();
let nshards = 1; let port = free_port();
let aof_path = dir.join("aof-0.aof");
with_runtime_configured(
port,
&dir,
nshards,
|rt| rt.with_auto_aof_rewrite(50, 16 * 1024),
|p| {
let mut c = std::net::TcpStream::connect(("127.0.0.1", p)).unwrap();
for rev in 0..800u32 {
c.write_all(&req(&[
b"SET",
b"counter",
format!("revision-number-padding-{rev:08}").as_bytes(),
]))
.unwrap();
read_reply(&mut c, b"+OK\r\n");
}
let post = wait_for_size_below_heartbeat(&aof_path, &mut c, 8 * 1024, 20_000);
assert!(
post < 8 * 1024,
"auto AOF rewrite did not fire: {post} bytes still on disk after \
800 SETs (un-rewritten would be ≈ 48 KiB)"
);
},
);
let port2 = free_port();
with_runtime(port2, &dir, nshards, |p| {
let mut c = std::net::TcpStream::connect(("127.0.0.1", p)).unwrap();
c.write_all(&req(&[b"GET", b"counter"])).unwrap();
read_reply(
&mut c,
b"$32\r\nrevision-number-padding-00000799\r\n",
);
});
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn auto_aof_rewrite_respects_pct_zero_disable() {
let dir = std::env::temp_dir().join(format!(
"kevy-auto-rewrite-off-{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
std::fs::create_dir_all(&dir).unwrap();
let nshards = 1;
let port = free_port();
let aof_path = dir.join("aof-0.aof");
with_runtime_configured(
port,
&dir,
nshards,
|rt| rt.with_auto_aof_rewrite(0, 1024),
|p| {
let mut c = std::net::TcpStream::connect(("127.0.0.1", p)).unwrap();
for rev in 0..800u32 {
c.write_all(&req(&[
b"SET",
b"k",
format!("revision-number-padding-{rev:08}").as_bytes(),
]))
.unwrap();
read_reply(&mut c, b"+OK\r\n");
}
let pre = wait_for_size_at_least_heartbeat(&aof_path, &mut c, 16 * 1024, 5_000);
assert!(pre >= 16 * 1024, "AOF did not grow: {pre} bytes");
let deadline = std::time::Instant::now() + std::time::Duration::from_millis(600);
while std::time::Instant::now() < deadline {
c.write_all(&req(&[b"PING"])).unwrap();
read_reply(&mut c, b"+PONG\r\n");
std::thread::sleep(std::time::Duration::from_millis(10));
}
let post = std::fs::metadata(&aof_path).map_or(0, |m| m.len());
assert!(
post >= pre,
"auto-rewrite fired despite pct=0: {post} vs {pre} pre"
);
},
);
let _ = std::fs::remove_dir_all(&dir);
}
fn wait_for_size_at_least_heartbeat(
path: &std::path::Path,
c: &mut std::net::TcpStream,
floor: u64,
timeout_ms: u64,
) -> u64 {
let deadline = std::time::Instant::now() + std::time::Duration::from_millis(timeout_ms);
loop {
let sz = std::fs::metadata(path).map_or(0, |m| m.len());
if sz >= floor || std::time::Instant::now() >= deadline {
return sz;
}
let _ = c.write_all(&req(&[b"PING"]));
read_reply(c, b"+PONG\r\n");
std::thread::sleep(std::time::Duration::from_millis(10));
}
}
fn wait_for_size_below_heartbeat(
path: &std::path::Path,
c: &mut std::net::TcpStream,
pre: u64,
timeout_ms: u64,
) -> u64 {
let deadline = std::time::Instant::now() + std::time::Duration::from_millis(timeout_ms);
loop {
let sz = std::fs::metadata(path).map_or(0, |m| m.len());
if sz < pre || std::time::Instant::now() >= deadline {
return sz;
}
let _ = c.write_all(&req(&[b"PING"]));
read_reply(c, b"+PONG\r\n");
std::thread::sleep(std::time::Duration::from_millis(10));
}
}
fn read_integer(s: &mut std::net::TcpStream) -> i64 {
let mut byte = [0u8; 1];
s.read_exact(&mut byte).unwrap();
assert_eq!(byte[0], b':', "expected RESP integer");
let mut n = Vec::new();
loop {
s.read_exact(&mut byte).unwrap();
if byte[0] == b'\r' {
s.read_exact(&mut byte).unwrap(); break;
}
n.push(byte[0]);
}
String::from_utf8(n).unwrap().parse().unwrap()
}
#[test]
fn relative_ttl_survives_restart_at_original_deadline() {
let dir = std::env::temp_dir().join(format!(
"kevy-ttl-restart-{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
std::fs::create_dir_all(&dir).unwrap();
let nshards = 2;
let port = free_port();
with_runtime(port, &dir, nshards, |p| {
let mut c = std::net::TcpStream::connect(("127.0.0.1", p)).unwrap();
c.write_all(&req(&[b"SET", b"k", b"v"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"PEXPIRE", b"k", b"100000"])).unwrap();
read_reply(&mut c, b":1\r\n");
});
std::thread::sleep(std::time::Duration::from_secs(3));
let port2 = free_port();
with_runtime(port2, &dir, nshards, |p| {
let mut c = std::net::TcpStream::connect(("127.0.0.1", p)).unwrap();
c.write_all(&req(&[b"GET", b"k"])).unwrap();
read_reply(&mut c, b"$1\r\nv\r\n"); c.write_all(&req(&[b"PTTL", b"k"])).unwrap();
let pttl = read_integer(&mut c);
assert!(
(0..=98_000).contains(&pttl),
"PTTL after restart = {pttl} ms; expected the original deadline \
(~97 s) minus downtime, not a reset to the full 100 s"
);
assert!(pttl > 90_000, "PTTL {pttl} ms implausibly low — key nearly gone");
});
let _ = std::fs::remove_dir_all(&dir);
}
fn build_grouped_stream(c: &mut std::net::TcpStream) {
let entry = |id: &str| format!("*2\r\n$3\r\n{id}\r\n*2\r\n$1\r\nf\r\n$1\r\nv\r\n");
for id in ["1-1", "2-1", "3-1"] {
c.write_all(&req(&[b"XADD", b"st", id.as_bytes(), b"f", b"v"])).unwrap();
read_reply(c, format!("$3\r\n{id}\r\n").as_bytes());
}
c.write_all(&req(&[b"XGROUP", b"CREATE", b"st", b"g", b"0"])).unwrap();
read_reply(c, b"+OK\r\n");
c.write_all(&req(&[
b"XREADGROUP", b"GROUP", b"g", b"c1", b"COUNT", b"2", b"STREAMS", b"st", b">",
]))
.unwrap();
read_reply(
c,
format!("*1\r\n*2\r\n$2\r\nst\r\n*2\r\n{}{}", entry("1-1"), entry("2-1")).as_bytes(),
);
c.write_all(&req(&[b"XREADGROUP", b"GROUP", b"g", b"c2", b"STREAMS", b"st", b">"]))
.unwrap();
read_reply(
c,
format!("*1\r\n*2\r\n$2\r\nst\r\n*1\r\n{}", entry("3-1")).as_bytes(),
);
c.write_all(&req(&[b"XDEL", b"st", b"2-1"])).unwrap();
read_reply(c, b":1\r\n");
c.write_all(&req(&[b"XADD", b"st2", b"5-1", b"f", b"v"])).unwrap();
read_reply(c, b"$3\r\n5-1\r\n");
c.write_all(&req(&[b"XDEL", b"st2", b"5-1"])).unwrap();
read_reply(c, b":1\r\n");
c.write_all(&req(&[b"XGROUP", b"CREATE", b"st2", b"g2", b"5-1"])).unwrap();
read_reply(c, b"+OK\r\n");
}
fn assert_grouped_stream_restored(c: &mut std::net::TcpStream, tombstone_kept: bool) {
let (total, c1) = if tombstone_kept { (3, 2) } else { (2, 1) };
c.write_all(&req(&[b"XPENDING", b"st", b"g"])).unwrap();
read_reply(
c,
format!(
"*4\r\n:{total}\r\n$3\r\n1-1\r\n$3\r\n3-1\r\n*2\r\n*2\r\n$2\r\nc1\r\n$1\r\n{c1}\r\n*2\r\n$2\r\nc2\r\n$1\r\n1\r\n"
)
.as_bytes(),
);
c.write_all(&req(&[b"XREADGROUP", b"GROUP", b"g", b"c1", b"STREAMS", b"st", b"0"]))
.unwrap();
read_reply(
c,
b"*1\r\n*2\r\n$2\r\nst\r\n*1\r\n*2\r\n$3\r\n1-1\r\n*2\r\n$1\r\nf\r\n$1\r\nv\r\n",
);
c.write_all(&req(&[b"XADD", b"st2", b"5-1", b"f", b"v"])).unwrap();
read_reply(
c,
b"-ERR The ID specified in XADD is equal or smaller than the target stream top item\r\n",
);
c.write_all(&req(&[b"XPENDING", b"st2", b"g2"])).unwrap();
read_reply(c, b"*4\r\n:0\r\n$-1\r\n$-1\r\n*-1\r\n");
}
#[test]
fn stream_groups_survive_bgrewriteaof_restart() {
let dir = std::env::temp_dir().join(format!(
"kevy-groups-aof-{}",
std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_nanos()
));
std::fs::create_dir_all(&dir).unwrap();
let port = free_port();
with_runtime(port, &dir, 1, |p| {
let mut c = std::net::TcpStream::connect(("127.0.0.1", p)).unwrap();
build_grouped_stream(&mut c);
c.write_all(&req(&[b"BGREWRITEAOF"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
wait_for("rewritten AOF to swap in", || {
std::fs::read(dir.join("aof-0.aof")).is_ok_and(|now| {
now.windows(6).any(|w| w == b"XCLAIM")
&& !now.windows(10).any(|w| w == b"XREADGROUP")
})
});
});
let port2 = free_port();
with_runtime(port2, &dir, 1, |p| {
let mut c = std::net::TcpStream::connect(("127.0.0.1", p)).unwrap();
assert_grouped_stream_restored(&mut c, false);
});
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn stream_groups_survive_save_restart() {
let dir = std::env::temp_dir().join(format!(
"kevy-groups-save-{}",
std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_nanos()
));
std::fs::create_dir_all(&dir).unwrap();
let port = free_port();
with_runtime(port, &dir, 1, |p| {
let mut c = std::net::TcpStream::connect(("127.0.0.1", p)).unwrap();
build_grouped_stream(&mut c);
c.write_all(&req(&[b"SAVE"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
});
let port2 = free_port();
with_runtime(port2, &dir, 1, |p| {
let mut c = std::net::TcpStream::connect(("127.0.0.1", p)).unwrap();
assert_grouped_stream_restored(&mut c, true);
});
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn bgsave_writes_snapshot_in_background_and_keeps_post_save_writes() {
let dir = std::env::temp_dir().join(format!(
"kevy-bgsave-{}",
std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_nanos()
));
std::fs::create_dir_all(&dir).unwrap();
let nshards = 4;
let port = free_port();
with_runtime(port, &dir, nshards, |p| {
let mut c = std::net::TcpStream::connect(("127.0.0.1", p)).unwrap();
for i in 0..50u32 {
c.write_all(&req(&[b"SET", format!("k{i}").as_bytes(), b"v"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
}
c.write_all(&req(&[b"BGSAVE"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
for i in 50..60u32 {
c.write_all(&req(&[b"SET", format!("k{i}").as_bytes(), b"v"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
}
wait_for("background snapshots to land", || {
(0..nshards).all(|s| dir.join(format!("dump-{s}.rdb")).exists())
});
wait_for("aof reset to swap in", || {
(0..nshards).all(|s| {
std::fs::read(dir.join(format!("aof-{s}.aof")))
.is_ok_and(|b| !b.windows(4).any(|w| w == b"\nk0\r".as_slice()))
})
});
});
let port2 = free_port();
with_runtime(port2, &dir, nshards, |p| {
let mut c = std::net::TcpStream::connect(("127.0.0.1", p)).unwrap();
for i in 0..60u32 {
c.write_all(&req(&[b"GET", format!("k{i}").as_bytes()])).unwrap();
read_reply(&mut c, b"$1\r\nv\r\n");
}
c.write_all(&req(&[b"DBSIZE"])).unwrap();
read_reply(&mut c, b":60\r\n");
});
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn info_persistence_reports_rewrite_completion() {
let dir = std::env::temp_dir().join(format!(
"kevy-info-persist-{}",
std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_nanos()
));
std::fs::create_dir_all(&dir).unwrap();
let port = free_port();
with_runtime(port, &dir, 1, |p| {
let mut c = std::net::TcpStream::connect(("127.0.0.1", p)).unwrap();
for i in 0..100u32 {
c.write_all(&req(&[b"SET", format!("k{i}").as_bytes(), b"v"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
}
c.write_all(&req(&[b"BGREWRITEAOF"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
let info = |c: &mut std::net::TcpStream| -> String {
c.write_all(&req(&[b"INFO", b"persistence"])).unwrap();
let mut one = [0u8; 1];
let mut hdr = Vec::new();
loop {
c.read_exact(&mut one).unwrap();
hdr.push(one[0]);
if hdr.ends_with(b"\r\n") {
break;
}
}
let len: usize =
String::from_utf8_lossy(&hdr[1..hdr.len() - 2]).parse().unwrap();
let mut body = vec![0u8; len + 2];
c.read_exact(&mut body).unwrap();
String::from_utf8_lossy(&body).into_owned()
};
wait_for("INFO to report the completed rewrite", || {
let s = info(&mut c);
s.contains("aof_rewrites_total:1") && s.contains("aof_rewrite_in_progress:0")
});
wait_for("INFO to report the AOF format", || info(&mut c).contains("aof_format:v2"));
});
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn save_does_not_block_reactor_for_disk_write() {
let dir = std::env::temp_dir().join(format!(
"kevy-save-async-{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
std::fs::create_dir_all(&dir).unwrap();
let nshards = 4;
let port = free_port();
with_runtime(port, &dir, nshards, |p| {
let mut c = std::net::TcpStream::connect(("127.0.0.1", p)).unwrap();
let big = vec![b'x'; 256];
for i in 0..20_000u32 {
let mut argv = req(&[b"SET", format!("k{i}").as_bytes(), &big]);
argv.extend_from_slice(&[]);
c.write_all(&argv).unwrap();
read_reply(&mut c, b"+OK\r\n");
}
let save_t0 = std::time::Instant::now();
c.write_all(&req(&[b"SAVE"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
let save_reply_us = save_t0.elapsed().as_micros();
assert!(
save_reply_us < 50_000,
"SAVE +OK took {save_reply_us} µs — expected <50 ms (\
reactor blocked? sync save regression?)"
);
let get_t0 = std::time::Instant::now();
c.write_all(&req(&[b"GET", b"k1"])).unwrap();
let mut prefix = [0u8; 7];
c.read_exact(&mut prefix).unwrap();
assert_eq!(&prefix, b"$256\r\nx");
let mut rest = vec![0u8; 255 + 2];
c.read_exact(&mut rest).unwrap();
let get_us = get_t0.elapsed().as_micros();
assert!(
get_us < 50_000,
"GET after SAVE took {get_us} µs — expected <50 ms (reactor blocked?)"
);
wait_for("background SAVE to land all shard dumps", || {
(0..nshards).all(|s| dir.join(format!("dump-{s}.rdb")).exists())
});
});
let port2 = free_port();
with_runtime(port2, &dir, nshards, |p| {
let mut c = std::net::TcpStream::connect(("127.0.0.1", p)).unwrap();
c.write_all(&req(&[b"DBSIZE"])).unwrap();
let mut buf = [0u8; 16];
let n = c.read(&mut buf).unwrap();
let reply = String::from_utf8_lossy(&buf[..n]).to_string();
assert!(
reply.starts_with(":20000\r\n"),
"DBSIZE after restart = {reply:?} (expected :20000\\r\\n — \
async SAVE failed to land before shutdown?)"
);
});
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn save_at_shutdown_drains_to_disk() {
let dir = std::env::temp_dir().join(format!(
"kevy-save-shutdown-{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
std::fs::create_dir_all(&dir).unwrap();
let nshards = 4;
let port = free_port();
with_runtime(port, &dir, nshards, |p| {
let mut c = std::net::TcpStream::connect(("127.0.0.1", p)).unwrap();
let big = vec![b'y'; 1024];
for i in 0..5_000u32 {
c.write_all(&req(&[b"SET", format!("k{i}").as_bytes(), &big]))
.unwrap();
read_reply(&mut c, b"+OK\r\n");
}
c.write_all(&req(&[b"SAVE"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
});
let dumps_after = (0..nshards)
.filter(|i| dir.join(format!("dump-{i}.rdb")).exists())
.count();
assert_eq!(
dumps_after, nshards,
"shutdown drain did not flush all shards' snapshots: \
only {dumps_after}/{nshards} dump-N.rdb files exist"
);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn relative_ttl_frames_do_not_reanchor_on_replay() {
let dir = std::env::temp_dir().join(format!(
"kevy-ttl-reanchor-{}",
std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_nanos()
));
std::fs::create_dir_all(&dir).unwrap();
let port = free_port();
with_runtime(port, &dir, 1, |p| {
let mut c = std::net::TcpStream::connect(("127.0.0.1", p)).unwrap();
c.write_all(&req(&[b"SETEX", b"grey", b"100", b"v"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"EXPIRE", b"grey2", b"100"])).unwrap(); let mut buf = [0u8; 64];
let _ = c.read(&mut buf).unwrap();
});
std::thread::sleep(std::time::Duration::from_millis(2500));
let port = free_port();
with_runtime(port, &dir, 1, |p| {
let mut c = std::net::TcpStream::connect(("127.0.0.1", p)).unwrap();
c.write_all(&req(&[b"PTTL", b"grey"])).unwrap();
let mut buf = [0u8; 64];
let n = c.read(&mut buf).unwrap();
let s = String::from_utf8_lossy(&buf[..n]);
let ttl: i64 = s.trim_start_matches(':').trim().parse().expect("integer PTTL");
assert!(ttl > 0, "key survived the restart: {s}");
assert!(
ttl <= 100_000 - 2_000,
"TTL re-anchored on replay: read {ttl}ms of an original 100000ms \
after >=2.5s elapsed — the AOF frame must carry an absolute deadline"
);
});
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn mset_and_rename_survive_a_restart() {
let dir = std::env::temp_dir().join(format!(
"kevy-persist-replayverbs-{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
std::fs::create_dir_all(&dir).unwrap();
let nshards = 1;
with_runtime(free_port(), &dir, nshards, |p| {
let mut c = std::net::TcpStream::connect(("127.0.0.1", p)).unwrap();
c.write_all(&req(&[b"CONFIG", b"SET", b"appendfsync", b"always"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"MSET", b"m:1", b"a", b"m:2", b"b"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"SET", b"src", b"v"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"RENAME", b"src", b"dst"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
});
with_runtime(free_port(), &dir, nshards, |p| {
let mut c = std::net::TcpStream::connect(("127.0.0.1", p)).unwrap();
c.write_all(&req(&[b"GET", b"m:1"])).unwrap();
read_reply(&mut c, b"$1\r\na\r\n");
c.write_all(&req(&[b"GET", b"m:2"])).unwrap();
read_reply(&mut c, b"$1\r\nb\r\n");
c.write_all(&req(&[b"GET", b"dst"])).unwrap();
read_reply(&mut c, b"$1\r\nv\r\n");
c.write_all(&req(&[b"GET", b"src"])).unwrap();
read_reply(&mut c, b"$-1\r\n");
});
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn cross_shard_rename_survives_a_restart() {
let dir = std::env::temp_dir().join(format!(
"kevy-persist-xrename-{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
std::fs::create_dir_all(&dir).unwrap();
let nshards = 4;
let cross = |a: &[u8], b: &[u8]| {
kevy_rt::shard_of_key(a, nshards, false) != kevy_rt::shard_of_key(b, nshards, false)
};
assert!(cross(b"src", b"dst"), "test fixture must be cross-shard");
assert!(cross(b"h:src", b"h:dst"), "hash fixture must be cross-shard");
assert!(cross(b"keep", b"taken"), "refusal fixture must be cross-shard");
with_runtime(free_port(), &dir, nshards, |p| {
let mut c = std::net::TcpStream::connect(("127.0.0.1", p)).unwrap();
c.write_all(&req(&[b"CONFIG", b"SET", b"appendfsync", b"always"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"SET", b"src", b"v"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"EXPIRE", b"src", b"1000"])).unwrap();
read_reply(&mut c, b":1\r\n");
c.write_all(&req(&[b"HSET", b"h:src", b"f", b"1"])).unwrap();
read_reply(&mut c, b":1\r\n");
c.write_all(&req(&[b"RENAME", b"src", b"dst"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"RENAME", b"h:src", b"h:dst"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"SET", b"keep", b"mine"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"SET", b"taken", b"theirs"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"RENAMENX", b"keep", b"taken"])).unwrap();
read_reply(&mut c, b":0\r\n");
});
with_runtime(free_port(), &dir, nshards, |p| {
let mut c = std::net::TcpStream::connect(("127.0.0.1", p)).unwrap();
c.write_all(&req(&[b"GET", b"dst"])).unwrap();
read_reply(&mut c, b"$1\r\nv\r\n");
c.write_all(&req(&[b"GET", b"src"])).unwrap();
read_reply(&mut c, b"$-1\r\n");
c.write_all(&req(&[b"TTL", b"dst"])).unwrap();
let mut buf = [0u8; 32];
let n = c.read(&mut buf).unwrap();
let ttl: i64 = String::from_utf8_lossy(&buf[1..n - 2]).parse().unwrap();
assert!((900..=1000).contains(&ttl), "TTL must survive the move: {ttl}");
c.write_all(&req(&[b"HGET", b"h:dst", b"f"])).unwrap();
read_reply(&mut c, b"$1\r\n1\r\n");
c.write_all(&req(&[b"EXISTS", b"h:src"])).unwrap();
read_reply(&mut c, b":0\r\n");
c.write_all(&req(&[b"GET", b"keep"])).unwrap();
read_reply(&mut c, b"$4\r\nmine\r\n");
c.write_all(&req(&[b"GET", b"taken"])).unwrap();
read_reply(&mut c, b"$6\r\ntheirs\r\n");
});
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn every_value_type_round_trips_through_a_snapshot() {
let dir = std::env::temp_dir().join(format!(
"kevy-persist-snaptypes-{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
std::fs::create_dir_all(&dir).unwrap();
let nshards = 2;
with_runtime_configured(free_port(), &dir, nshards, |rt| rt.with_aof(false), |p| {
let mut c = std::net::TcpStream::connect(("127.0.0.1", p)).unwrap();
c.write_all(&req(&[b"SET", b"str", b"v"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"EXPIRE", b"str", b"500"])).unwrap();
read_reply(&mut c, b":1\r\n");
c.write_all(&req(&[b"HSET", b"h", b"f1", b"1", b"f2", b"2"])).unwrap();
read_reply(&mut c, b":2\r\n");
c.write_all(&req(&[b"RPUSH", b"l", b"a", b"b", b"c"])).unwrap();
read_reply(&mut c, b":3\r\n");
c.write_all(&req(&[b"SADD", b"s", b"m1", b"m2"])).unwrap();
read_reply(&mut c, b":2\r\n");
c.write_all(&req(&[b"ZADD", b"z", b"2.5", b"zn"])).unwrap();
read_reply(&mut c, b":1\r\n");
c.write_all(&req(&[b"HSET", b"hx", b"g", b"1"])).unwrap();
read_reply(&mut c, b":1\r\n");
c.write_all(&req(&[b"HEXPIRE", b"hx", b"400", b"FIELDS", b"1", b"g"])).unwrap();
read_reply(&mut c, b"*1\r\n:1\r\n");
c.write_all(&req(&[b"XADD", b"strm", b"1-1", b"f", b"v"])).unwrap();
read_reply(&mut c, b"$3\r\n1-1\r\n");
c.write_all(&req(&[b"SAVE"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
});
let dumps = (0..nshards).filter(|i| dir.join(format!("dump-{i}.rdb")).exists()).count();
assert!(dumps > 0, "no snapshot was written — nothing to round-trip");
let aofs = std::fs::read_dir(&dir)
.unwrap()
.filter_map(Result::ok)
.filter(|e| e.file_name().to_string_lossy().ends_with(".aof"))
.count();
assert_eq!(aofs, 0, "an AOF exists, so the snapshot is not what is under test");
with_runtime_configured(free_port(), &dir, nshards, |rt| rt.with_aof(false), |p| {
let mut c = std::net::TcpStream::connect(("127.0.0.1", p)).unwrap();
c.write_all(&req(&[b"GET", b"str"])).unwrap();
read_reply(&mut c, b"$1\r\nv\r\n");
c.write_all(&req(&[b"HGET", b"h", b"f2"])).unwrap();
read_reply(&mut c, b"$1\r\n2\r\n");
c.write_all(&req(&[b"LRANGE", b"l", b"0", b"-1"])).unwrap();
read_reply(&mut c, b"*3\r\n$1\r\na\r\n$1\r\nb\r\n$1\r\nc\r\n");
c.write_all(&req(&[b"SISMEMBER", b"s", b"m2"])).unwrap();
read_reply(&mut c, b":1\r\n");
c.write_all(&req(&[b"ZSCORE", b"z", b"zn"])).unwrap();
read_reply(&mut c, b"$3\r\n2.5\r\n");
c.write_all(&req(&[b"XRANGE", b"strm", b"-", b"+"])).unwrap();
read_reply(&mut c, b"*1\r\n*2\r\n$3\r\n1-1\r\n*2\r\n$1\r\nf\r\n$1\r\nv\r\n");
for (probe, floor) in [
(req(&[b"TTL", b"str"]), 400i64),
(req(&[b"HTTL", b"hx", b"FIELDS", b"1", b"g"]), 300i64),
] {
c.write_all(&probe).unwrap();
let mut buf = [0u8; 64];
let n = c.read(&mut buf).unwrap();
let text = String::from_utf8_lossy(&buf[..n]).into_owned();
let secs: i64 = text
.rsplit(':')
.next()
.and_then(|t| t.trim_end_matches("\r\n").parse().ok())
.unwrap_or(-1);
assert!(secs > floor, "TTL must survive the snapshot, got {text:?}");
}
});
let _ = std::fs::remove_dir_all(&dir);
}