use std::io::{Read, Write};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};
static START_GATE: Mutex<()> = Mutex::new(());
fn free_port() -> u16 {
let l = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
l.local_addr().unwrap().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 {:?}, got {:?}",
String::from_utf8_lossy(expected),
String::from_utf8_lossy(&buf),
);
}
struct Server {
port: u16,
dir: std::path::PathBuf,
stop: Arc<AtomicBool>,
handle: Option<std::thread::JoinHandle<()>>,
}
impl Server {
fn start(nshards: usize) -> Server {
let _gate = START_GATE.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
let port = free_port();
let dir = std::env::temp_dir().join(format!(
"kevy-watch-{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
std::fs::create_dir_all(&dir).unwrap();
let stop = Arc::new(AtomicBool::new(false));
let stop_thread = stop.clone();
let dir_thread = dir.clone();
let handle = std::thread::spawn(move || {
let rt = kevy_rt::Runtime::new([127, 0, 0, 1], port, nshards, kevy::KevyCommands)
.with_data_dir(dir_thread);
rt.run(stop_thread).unwrap();
});
let mut ready = false;
for _ in 0..200 {
if std::net::TcpStream::connect(("127.0.0.1", port)).is_ok() {
ready = true;
break;
}
std::thread::sleep(std::time::Duration::from_millis(5));
}
assert!(ready, "runtime did not come up");
Server {
port,
dir,
stop,
handle: Some(handle),
}
}
fn connect(&self) -> std::net::TcpStream {
let s = std::net::TcpStream::connect(("127.0.0.1", self.port)).unwrap();
s.set_read_timeout(Some(std::time::Duration::from_secs(10)))
.unwrap();
s
}
}
impl Drop for Server {
fn drop(&mut self) {
self.stop.store(true, Ordering::Relaxed);
if let Some(h) = self.handle.take() {
let _ = h.join();
}
let _ = std::fs::remove_dir_all(&self.dir);
}
}
#[test]
fn watch_then_exec_clean_commits() {
let srv = Server::start(1);
let mut c = srv.connect();
c.write_all(&req(&[b"WATCH", b"k"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"MULTI"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"SET", b"k", b"v"])).unwrap();
read_reply(&mut c, b"+QUEUED\r\n");
c.write_all(&req(&[b"EXEC"])).unwrap();
read_reply(&mut c, b"*1\r\n+OK\r\n");
c.write_all(&req(&[b"GET", b"k"])).unwrap();
read_reply(&mut c, b"$1\r\nv\r\n");
}
#[test]
fn watch_then_concurrent_write_aborts() {
let srv = Server::start(1);
let mut c = srv.connect();
let mut other = srv.connect();
c.write_all(&req(&[b"SET", b"k", b"orig"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"WATCH", b"k"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
other.write_all(&req(&[b"SET", b"k", b"stomp"])).unwrap();
read_reply(&mut other, b"+OK\r\n");
c.write_all(&req(&[b"MULTI"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"SET", b"k", b"v"])).unwrap();
read_reply(&mut c, b"+QUEUED\r\n");
c.write_all(&req(&[b"EXEC"])).unwrap();
read_reply(&mut c, b"*-1\r\n");
c.write_all(&req(&[b"GET", b"k"])).unwrap();
read_reply(&mut c, b"$5\r\nstomp\r\n");
}
#[test]
fn unwatch_clears_watched_set() {
let srv = Server::start(1);
let mut c = srv.connect();
let mut other = srv.connect();
c.write_all(&req(&[b"WATCH", b"k"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"UNWATCH"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
other.write_all(&req(&[b"SET", b"k", b"stomp"])).unwrap();
read_reply(&mut other, b"+OK\r\n");
c.write_all(&req(&[b"MULTI"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"SET", b"k", b"v"])).unwrap();
read_reply(&mut c, b"+QUEUED\r\n");
c.write_all(&req(&[b"EXEC"])).unwrap();
read_reply(&mut c, b"*1\r\n+OK\r\n");
c.write_all(&req(&[b"GET", b"k"])).unwrap();
read_reply(&mut c, b"$1\r\nv\r\n");
}
#[test]
fn watch_inside_multi_is_an_error() {
let srv = Server::start(1);
let mut c = srv.connect();
c.write_all(&req(&[b"MULTI"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"WATCH", b"k"])).unwrap();
let mut buf = [0u8; 64];
let n = c.read(&mut buf).unwrap();
assert!(
buf[..n].starts_with(b"-ERR WATCH inside MULTI"),
"got {:?}",
String::from_utf8_lossy(&buf[..n])
);
c.write_all(&req(&[b"DISCARD"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
}
#[test]
fn discard_clears_watched_set() {
let srv = Server::start(1);
let mut c = srv.connect();
let mut other = srv.connect();
c.write_all(&req(&[b"WATCH", b"k"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"MULTI"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"DISCARD"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
other.write_all(&req(&[b"SET", b"k", b"stomp"])).unwrap();
read_reply(&mut other, b"+OK\r\n");
c.write_all(&req(&[b"MULTI"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"SET", b"k", b"v"])).unwrap();
read_reply(&mut c, b"+QUEUED\r\n");
c.write_all(&req(&[b"EXEC"])).unwrap();
read_reply(&mut c, b"*1\r\n+OK\r\n");
}
#[test]
fn watch_no_stomp_then_exec_commits() {
let srv = Server::start(1);
let mut c = srv.connect();
c.write_all(&req(&[b"WATCH", b"fresh"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"MULTI"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"SET", b"fresh", b"x"])).unwrap();
read_reply(&mut c, b"+QUEUED\r\n");
c.write_all(&req(&[b"EXEC"])).unwrap();
read_reply(&mut c, b"*1\r\n+OK\r\n");
}
#[test]
fn watch_many_keys_cross_shard_clean_commits() {
let srv = Server::start(4);
let mut c = srv.connect();
let mut watch_req = vec![b"WATCH".as_slice()];
let keys: Vec<String> = (0..24).map(|i| format!("xs:{i}")).collect();
for k in &keys {
watch_req.push(k.as_bytes());
}
c.write_all(&req(&watch_req)).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"MULTI"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"SET", b"unrelated", b"v"])).unwrap();
read_reply(&mut c, b"+QUEUED\r\n");
c.write_all(&req(&[b"EXEC"])).unwrap();
read_reply(&mut c, b"*1\r\n+OK\r\n");
}
#[test]
fn watch_many_keys_cross_shard_stomp_aborts() {
let srv = Server::start(4);
let mut c = srv.connect();
let mut other = srv.connect();
let mut watch_req = vec![b"WATCH".as_slice()];
let keys: Vec<String> = (0..24).map(|i| format!("xs:{i}")).collect();
for k in &keys {
watch_req.push(k.as_bytes());
}
c.write_all(&req(&watch_req)).unwrap();
read_reply(&mut c, b"+OK\r\n");
other.write_all(&req(&[b"SET", b"xs:12", b"stomp"])).unwrap();
read_reply(&mut other, b"+OK\r\n");
c.write_all(&req(&[b"MULTI"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"SET", b"unrelated", b"v"])).unwrap();
read_reply(&mut c, b"+QUEUED\r\n");
c.write_all(&req(&[b"EXEC"])).unwrap();
read_reply(&mut c, b"*-1\r\n");
}
#[test]
fn exec_after_watch_with_multi_cmd_queue() {
let srv = Server::start(4);
let mut c = srv.connect();
c.write_all(&req(&[b"WATCH", b"q:a"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"MULTI"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"SET", b"q:a", b"0"])).unwrap();
read_reply(&mut c, b"+QUEUED\r\n");
c.write_all(&req(&[b"INCR", b"q:a"])).unwrap();
read_reply(&mut c, b"+QUEUED\r\n");
c.write_all(&req(&[b"PING"])).unwrap();
read_reply(&mut c, b"+QUEUED\r\n");
c.write_all(&req(&[b"GET", b"q:a"])).unwrap();
read_reply(&mut c, b"+QUEUED\r\n");
c.write_all(&req(&[b"EXEC"])).unwrap();
read_reply(
&mut c,
b"*4\r\n+OK\r\n:1\r\n+PONG\r\n$1\r\n1\r\n",
);
}
#[test]
fn watched_exec_with_queued_multi_target_del() {
let srv = Server::start(4);
let mut c = srv.connect();
for k in ["mk:a", "mk:b", "mk:c"] {
c.write_all(&req(&[b"SET", k.as_bytes(), b"v"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
}
c.write_all(&req(&[b"WATCH", b"sentinel"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"MULTI"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"DEL", b"mk:a", b"mk:b", b"mk:c"])).unwrap();
read_reply(&mut c, b"+QUEUED\r\n");
c.write_all(&req(&[b"EXISTS", b"mk:a", b"mk:b", b"mk:c"])).unwrap();
read_reply(&mut c, b"+QUEUED\r\n");
c.write_all(&req(&[b"EXEC"])).unwrap();
read_reply(&mut c, b"*2\r\n:3\r\n:0\r\n");
}
#[test]
fn watched_exec_with_queued_publish_emits_error_placeholder() {
let srv = Server::start(4);
let mut c = srv.connect();
c.write_all(&req(&[b"WATCH", b"pw:k"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"MULTI"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"PUBLISH", b"ch", b"x"])).unwrap();
read_reply(&mut c, b"+QUEUED\r\n");
c.write_all(&req(&[b"INCR", b"pw:k"])).unwrap();
read_reply(&mut c, b"+QUEUED\r\n");
c.write_all(&req(&[b"EXEC"])).unwrap();
let want = b"*2\r\n-ERR pub/sub or WATCH or HELLO or RENAME not allowed inside MULTI in v2-3a (queued-RENAME orchestration pending v2-3b)\r\n:1\r\n";
read_reply(&mut c, want);
}
#[test]
fn watched_exec_with_queued_unwatch_is_ok_noop() {
let srv = Server::start(4);
let mut c = srv.connect();
c.write_all(&req(&[b"WATCH", b"u:k"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"MULTI"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"UNWATCH"])).unwrap();
read_reply(&mut c, b"+QUEUED\r\n");
c.write_all(&req(&[b"INCR", b"u:k"])).unwrap();
read_reply(&mut c, b"+QUEUED\r\n");
c.write_all(&req(&[b"EXEC"])).unwrap();
read_reply(&mut c, b"*2\r\n+OK\r\n:1\r\n");
}
#[test]
fn watched_exec_quit_inside_multi_closes_after_drain() {
let srv = Server::start(4);
let mut c = srv.connect();
c.write_all(&req(&[b"WATCH", b"q:k"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"MULTI"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"SET", b"q:k", b"v"])).unwrap();
read_reply(&mut c, b"+QUEUED\r\n");
c.write_all(&req(&[b"QUIT"])).unwrap();
read_reply(&mut c, b"+QUEUED\r\n");
c.write_all(&req(&[b"EXEC"])).unwrap();
read_reply(&mut c, b"*2\r\n+OK\r\n+OK\r\n");
let mut buf = [0u8; 32];
let n = c.read(&mut buf).unwrap_or(0);
assert_eq!(n, 0, "expected EOF after queued QUIT inside EXEC");
}
#[test]
fn pipelined_command_after_exec_is_unaffected() {
let srv = Server::start(4);
let mut c = srv.connect();
let mut other = srv.connect();
c.write_all(&req(&[b"SET", b"px", b"orig"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
c.write_all(&req(&[b"WATCH", b"px"])).unwrap();
read_reply(&mut c, b"+OK\r\n");
other.write_all(&req(&[b"SET", b"px", b"stomp"])).unwrap();
read_reply(&mut other, b"+OK\r\n");
let mut batch = Vec::new();
batch.extend_from_slice(&req(&[b"MULTI"]));
batch.extend_from_slice(&req(&[b"SET", b"px", b"v"]));
batch.extend_from_slice(&req(&[b"EXEC"]));
batch.extend_from_slice(&req(&[b"GET", b"px"]));
c.write_all(&batch).unwrap();
read_reply(&mut c, b"+OK\r\n"); read_reply(&mut c, b"+QUEUED\r\n"); read_reply(&mut c, b"*-1\r\n"); read_reply(&mut c, b"$5\r\nstomp\r\n"); }