use std::process::{Child, Command};
use kevy_cli::migrate::{run_export, run_import};
use kevy_resp_client::RespClient;
struct Srv {
child: Child,
port: u16,
}
impl Srv {
fn start() -> Srv {
let port = std::net::TcpListener::bind("127.0.0.1:0").unwrap().local_addr().unwrap().port();
let bin = std::path::Path::new(env!("CARGO_BIN_EXE_kevy-cli"))
.parent()
.unwrap()
.join("kevy");
if !bin.exists() {
let cargo = std::env::var("CARGO").unwrap_or_else(|_| "cargo".into());
let status = Command::new(cargo)
.args(["build", "-p", "kevy", "--bin", "kevy"])
.status()
.expect("spawn cargo build");
assert!(status.success(), "cargo build -p kevy --bin kevy failed");
}
assert!(bin.exists(), "kevy server binary still missing at {bin:?}");
let child = Command::new(&bin)
.args(["--port", &port.to_string(), "--threads", "2", "--no-aof"])
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.spawn()
.expect("spawn kevy server");
for _ in 0..100 {
if std::net::TcpStream::connect(("127.0.0.1", port)).is_ok() {
break;
}
std::thread::sleep(std::time::Duration::from_millis(20));
}
Srv { child, port }
}
fn client(&self) -> RespClient {
RespClient::connect("127.0.0.1", self.port).unwrap()
}
}
impl Drop for Srv {
fn drop(&mut self) {
let _ = self.child.kill();
let _ = self.child.wait();
}
}
fn digest(c: &mut RespClient, prefix: &str) -> String {
let r = c.request_borrowed(&[b"PREFIX.DIGEST", prefix.as_bytes()]).unwrap();
format!("{r:?}")
}
#[test]
fn export_import_roundtrip_digest_equal() {
let src = Srv::start();
let dst = Srv::start();
let mut cs = src.client();
for i in 0..500 {
cs.request_borrowed(&[b"SET", format!("mig:s:{i}").as_bytes(), format!("v{i}").as_bytes()]).unwrap();
cs.request_borrowed(&[
b"HSET", format!("mig:h:{i}").as_bytes(), b"a", format!("{i}").as_bytes(), b"b", b"x",
]).unwrap();
}
cs.request_borrowed(&[b"RPUSH", b"mig:list", b"1", b"2", b"3"]).unwrap();
cs.request_borrowed(&[b"SADD", b"mig:set", b"p", b"q"]).unwrap();
cs.request_borrowed(&[b"ZADD", b"mig:zset", b"1.5", b"m", b"2.5", b"n"]).unwrap();
cs.request_borrowed(&[b"SET", b"mig:ttl", b"soon"]).unwrap();
cs.request_borrowed(&[b"PEXPIRE", b"mig:ttl", b"60000"]).unwrap();
cs.request_borrowed(&[b"SET", b"nomig:1", b"skip"]).unwrap();
let dir = std::env::temp_dir().join(format!("kevy-mig-{}", std::process::id()));
std::fs::create_dir_all(&dir).unwrap();
let file = dir.join("dump.resp");
let n = run_export(&mut cs, Some(b"mig:"), &file).unwrap();
assert_eq!(n, 1004, "500+500 rows + list/set/zset/ttl");
let mut cd = dst.client();
let rep = run_import(&mut cd, &file, false, true).unwrap();
assert_eq!(rep.errors, 0);
assert!(rep.sent >= 1004, "at least one frame per key: {}", rep.sent);
let ds = digest(&mut cs, "mig:");
let dd = digest(&mut cd, "mig:");
assert_eq!(ds, dd, "src {ds} vs dst {dd}");
let r = cd.request_borrowed(&[b"EXISTS", b"nomig:1"]).unwrap();
assert_eq!(format!("{r:?}"), "Int(0)");
let r = cd.request_borrowed(&[b"PTTL", b"mig:ttl"]).unwrap();
let s = format!("{r:?}");
let ms: i64 = s.trim_start_matches("Int(").trim_end_matches(')').parse().unwrap();
assert!(ms > 30_000 && ms <= 60_000, "absolute TTL carried: {ms}");
let file2 = dir.join("dump2.resp");
let n2 = run_export(&mut cs, Some(b"mig:"), &file2).unwrap();
assert_eq!(n2, n);
std::fs::write(file2.with_extension("progress"), b"0").unwrap();
let rep2 = run_import(&mut cd, &file2, true, true).unwrap();
assert_eq!(rep2.errors, 0, "idempotent replay");
assert_eq!(digest(&mut cd, "mig:"), ds);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn bulk_ops_and_diff() {
let srv = Srv::start();
let mut c = srv.client();
for i in 0..200 {
c.request_borrowed(&[b"SET", format!("bk:{i}").as_bytes(), format!("v{i}").as_bytes()]).unwrap();
}
c.request_borrowed(&[b"SET", b"bk:ttl", b"x"]).unwrap();
c.request_borrowed(&[b"PEXPIRE", b"bk:ttl", b"60000"]).unwrap();
let n = kevy_cli::bulk::run_copy_prefix(&mut c, b"bk:", b"ck:", 0).unwrap();
assert_eq!(n, 201);
let (na, da) = kevy_cli::bulk::run_digest(&mut c, b"bk:").unwrap();
let (nb, _db) = kevy_cli::bulk::run_digest(&mut c, b"ck:").unwrap();
assert_eq!((na, nb), (201, 201));
let r = c.request_borrowed(&[b"GET", b"ck:42"]).unwrap();
assert_eq!(format!("{r:?}"), format!("{:?}", kevy_cli::Reply::Bulk(b"v42".to_vec())));
let r = c.request_borrowed(&[b"PTTL", b"ck:ttl"]).unwrap();
let ms: i64 = format!("{r:?}").trim_start_matches("Int(").trim_end_matches(')').parse().unwrap();
assert!(ms > 0, "COPY carries TTL: {ms}");
let _ = da;
let n = kevy_cli::bulk::run_delete_prefix(&mut c, b"ck:", 0, true).unwrap();
assert_eq!(n, 201);
let (still, _) = kevy_cli::bulk::run_digest(&mut c, b"ck:").unwrap();
assert_eq!(still, 201);
let t0 = std::time::Instant::now();
let n = kevy_cli::bulk::run_delete_prefix(&mut c, b"ck:", 400, false).unwrap();
let dt = t0.elapsed().as_secs_f64();
assert_eq!(n, 201);
assert!(dt > 0.3, "rate limit engaged: {dt:.2}s");
let (gone, _) = kevy_cli::bulk::run_digest(&mut c, b"ck:").unwrap();
assert_eq!(gone, 0);
let srv2 = Srv::start();
let mut c2 = srv2.client();
let mut c1b = srv.client();
let mut out = Vec::new();
let bad = kevy_cli::bulk::run_diff(&mut c1b, &mut c2, &[b"bk:".to_vec()], &mut out).unwrap();
assert_eq!(bad.len(), 1, "empty dst mismatches: {}", String::from_utf8_lossy(&out));
let dir = std::env::temp_dir().join(format!("kevy-bulk-{}", std::process::id()));
std::fs::create_dir_all(&dir).unwrap();
let file = dir.join("bk.resp");
let mut c1c = srv.client();
kevy_cli::migrate::run_export(&mut c1c, Some(b"bk:"), &file).unwrap();
kevy_cli::migrate::run_import(&mut c2, &file, false, true).unwrap();
let mut out = Vec::new();
let bad = kevy_cli::bulk::run_diff(&mut c1c, &mut c2, &[b"bk:".to_vec()], &mut out).unwrap();
assert!(bad.is_empty(), "{}", String::from_utf8_lossy(&out));
let mut out = Vec::new();
kevy_cli::bulk::run_inspect(&mut c1c, b"bk:", &mut out).unwrap();
let s = String::from_utf8_lossy(&out);
assert!(s.contains("201 keys") && s.contains("string: 201"), "{s}");
let _ = std::fs::remove_dir_all(&dir);
}