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_n(s: &mut std::net::TcpStream, n: usize) -> Vec<u8> {
let mut buf = vec![0u8; n];
s.read_exact(&mut buf).unwrap();
buf
}
fn read_line(s: &mut std::net::TcpStream, out: &mut Vec<u8>) {
loop {
let b = read_n(s, 1);
out.extend_from_slice(&b);
if out.ends_with(b"\r\n") {
break;
}
}
}
fn read_len(s: &mut std::net::TcpStream, out: &mut Vec<u8>) -> i64 {
let start = out.len();
read_line(s, out);
let line = &out[start..out.len() - 2];
std::str::from_utf8(line).unwrap().parse().unwrap()
}
fn read_reply(s: &mut std::net::TcpStream) -> Vec<u8> {
let head = read_n(s, 1);
let mut out = head.clone();
match head[0] {
b'+' | b'-' | b':' => read_line(s, &mut out),
b'$' => {
let len = read_len(s, &mut out);
if len < 0 {
return out;
}
out.extend_from_slice(&read_n(s, len as usize + 2));
}
b'*' => {
let n = read_len(s, &mut out);
if n < 0 {
return out;
}
for _ in 0..n {
out.extend_from_slice(&read_reply(s));
}
}
other => panic!("unknown reply prefix {other:?}"),
}
out
}
struct Server {
port: u16,
dir: std::path::PathBuf,
stop: Arc<AtomicBool>,
handle: Option<std::thread::JoinHandle<()>>,
}
impl Server {
fn start(nshards: usize) -> Self {
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-blocking-{}",
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();
});
for _ in 0..200 {
if std::net::TcpStream::connect(("127.0.0.1", port)).is_ok() {
break;
}
std::thread::sleep(std::time::Duration::from_millis(5));
}
Self {
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(5)))
.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 blpop_returns_immediately_when_list_has_value() {
let srv = Server::start(1);
let mut c = srv.connect();
c.write_all(&req(&[b"RPUSH", b"k", b"v"])).unwrap();
let _ = read_reply(&mut c); c.write_all(&req(&[b"BLPOP", b"k", b"5"])).unwrap();
let reply = read_reply(&mut c);
assert_eq!(reply, b"*2\r\n$1\r\nk\r\n$1\r\nv\r\n");
}
#[test]
fn blpop_times_out_with_nil_array_when_list_empty() {
let srv = Server::start(1);
let mut c = srv.connect();
c.write_all(&req(&[b"BLPOP", b"empty", b"0.1"])).unwrap();
let t0 = std::time::Instant::now();
let reply = read_reply(&mut c);
let elapsed = t0.elapsed();
assert_eq!(reply, b"*-1\r\n", "BLPOP timeout must return nil array");
assert!(
elapsed >= std::time::Duration::from_millis(80),
"BLPOP should block at least the requested timeout (~100ms), got {elapsed:?}",
);
assert!(
elapsed < std::time::Duration::from_secs(2),
"BLPOP timeout fired far too late ({elapsed:?})",
);
}
#[test]
fn blpop_woken_by_concurrent_push() {
let srv = Server::start(1);
let mut consumer = srv.connect();
let mut producer = srv.connect();
consumer
.write_all(&req(&[b"BLPOP", b"wakeable", b"5"]))
.unwrap();
std::thread::sleep(std::time::Duration::from_millis(50));
producer
.write_all(&req(&[b"LPUSH", b"wakeable", b"hello"]))
.unwrap();
let _push_reply = read_reply(&mut producer); let reply = read_reply(&mut consumer);
assert_eq!(reply, b"*2\r\n$8\r\nwakeable\r\n$5\r\nhello\r\n");
}
#[test]
fn brpop_returns_immediately_when_list_has_value() {
let srv = Server::start(1);
let mut c = srv.connect();
c.write_all(&req(&[b"RPUSH", b"k", b"v"])).unwrap();
let _ = read_reply(&mut c);
c.write_all(&req(&[b"BRPOP", b"k", b"5"])).unwrap();
let reply = read_reply(&mut c);
assert_eq!(reply, b"*2\r\n$1\r\nk\r\n$1\r\nv\r\n");
}
#[test]
fn brpop_woken_by_concurrent_rpush() {
let srv = Server::start(1);
let mut consumer = srv.connect();
let mut producer = srv.connect();
consumer.write_all(&req(&[b"BRPOP", b"q", b"5"])).unwrap();
std::thread::sleep(std::time::Duration::from_millis(50));
producer.write_all(&req(&[b"RPUSH", b"q", b"x"])).unwrap();
let _ = read_reply(&mut producer);
let reply = read_reply(&mut consumer);
assert_eq!(reply, b"*2\r\n$1\r\nq\r\n$1\r\nx\r\n");
}
#[test]
fn xread_block_returns_immediately_when_stream_has_entry() {
let srv = Server::start(1);
let mut c = srv.connect();
c.write_all(&req(&[b"XADD", b"s", b"1-0", b"f", b"v"])).unwrap();
let _ = read_reply(&mut c); c.write_all(&req(&[
b"XREAD", b"BLOCK", b"5000", b"STREAMS", b"s", b"0",
])).unwrap();
let reply = read_reply(&mut c);
assert!(reply.starts_with(b"*1\r\n"), "expected one stream in reply, got {reply:?}");
assert!(reply.windows(3).any(|w| w == b"1-0"), "expected entry id 1-0 in reply");
assert!(reply.windows(1).any(|w| w == b"v"), "expected value v in reply");
}
#[test]
fn xread_block_times_out_with_nil_bulk_when_no_entries() {
let srv = Server::start(1);
let mut c = srv.connect();
c.write_all(&req(&[b"XADD", b"s", b"1-0", b"f", b"v"])).unwrap();
let _ = read_reply(&mut c);
c.write_all(&req(&[
b"XREAD", b"BLOCK", b"100", b"STREAMS", b"s", b"$",
])).unwrap();
let t0 = std::time::Instant::now();
let reply = read_reply(&mut c);
let elapsed = t0.elapsed();
assert_eq!(reply, b"$-1\r\n", "XREAD BLOCK timeout must return nil bulk");
assert!(elapsed >= std::time::Duration::from_millis(80));
}
#[test]
fn xread_block_woken_by_concurrent_xadd() {
let srv = Server::start(1);
let mut consumer = srv.connect();
let mut producer = srv.connect();
producer
.write_all(&req(&[b"XADD", b"stream", b"1-0", b"f", b"v"]))
.unwrap();
let _ = read_reply(&mut producer);
consumer
.write_all(&req(&[
b"XREAD", b"BLOCK", b"5000", b"STREAMS", b"stream", b"1-0",
]))
.unwrap();
std::thread::sleep(std::time::Duration::from_millis(50));
producer
.write_all(&req(&[b"XADD", b"stream", b"2-0", b"f", b"v2"]))
.unwrap();
let _ = read_reply(&mut producer);
let reply = read_reply(&mut consumer);
assert!(reply.starts_with(b"*1\r\n"));
assert!(reply.windows(3).any(|w| w == b"2-0"));
assert!(reply.windows(2).any(|w| w == b"v2"));
}
#[test]
fn xread_block_dollar_id_wakes() {
let srv = Server::start(1);
let mut consumer = srv.connect();
let mut producer = srv.connect();
producer
.write_all(&req(&[b"XADD", b"stream", b"1-0", b"f", b"v"]))
.unwrap();
let _ = read_reply(&mut producer);
consumer
.write_all(&req(&[
b"XREAD", b"BLOCK", b"5000", b"STREAMS", b"stream", b"$",
]))
.unwrap();
std::thread::sleep(std::time::Duration::from_millis(50));
producer
.write_all(&req(&[b"XADD", b"stream", b"2-0", b"f", b"v2"]))
.unwrap();
let _ = read_reply(&mut producer);
let reply = read_reply(&mut consumer);
assert!(
reply.starts_with(b"*1\r\n"),
"expected one stream in reply, got {:?}",
std::str::from_utf8(&reply).unwrap_or("<non-utf8>")
);
assert!(reply.windows(3).any(|w| w == b"2-0"));
assert!(reply.windows(2).any(|w| w == b"v2"));
}
#[test]
fn xreadgroup_block_times_out_when_no_new_entries() {
let srv = Server::start(1);
let mut c = srv.connect();
c.write_all(&req(&[b"XADD", b"s", b"1-0", b"f", b"v"])).unwrap();
let _ = read_reply(&mut c);
c.write_all(&req(&[b"XGROUP", b"CREATE", b"s", b"g", b"$"])).unwrap();
let _ = read_reply(&mut c); c.write_all(&req(&[
b"XREADGROUP",
b"GROUP",
b"g",
b"alice",
b"BLOCK",
b"100",
b"STREAMS",
b"s",
b">",
])).unwrap();
let t0 = std::time::Instant::now();
let reply = read_reply(&mut c);
let elapsed = t0.elapsed();
assert_eq!(reply, b"$-1\r\n", "XREADGROUP BLOCK timeout returns nil bulk");
assert!(elapsed >= std::time::Duration::from_millis(80));
}
#[test]
fn xreadgroup_block_woken_by_concurrent_xadd() {
let srv = Server::start(1);
let mut consumer = srv.connect();
let mut producer = srv.connect();
producer.write_all(&req(&[b"XADD", b"stream2", b"1-0", b"f", b"v"])).unwrap();
let _ = read_reply(&mut producer);
producer.write_all(&req(&[b"XGROUP", b"CREATE", b"stream2", b"g", b"$"])).unwrap();
let _ = read_reply(&mut producer);
consumer.write_all(&req(&[
b"XREADGROUP",
b"GROUP",
b"g",
b"bob",
b"BLOCK",
b"5000",
b"STREAMS",
b"stream2",
b">",
])).unwrap();
std::thread::sleep(std::time::Duration::from_millis(50));
producer.write_all(&req(&[b"XADD", b"stream2", b"2-0", b"f", b"v2"])).unwrap();
let _ = read_reply(&mut producer);
let reply = read_reply(&mut consumer);
assert!(reply.starts_with(b"*1\r\n"));
assert!(reply.windows(3).any(|w| w == b"2-0"));
assert!(reply.windows(2).any(|w| w == b"v2"));
}
#[test]
fn blpop_multi_key_times_out_when_all_empty() {
let srv = Server::start(1);
let mut c = srv.connect();
c.write_all(&req(&[b"BLPOP", b"a", b"b", b"0.1"])).unwrap();
let t0 = std::time::Instant::now();
let reply = read_reply(&mut c);
assert_eq!(reply, b"*-1\r\n", "multi-key BLPOP timeout returns nil array");
assert!(t0.elapsed() >= std::time::Duration::from_millis(80));
}
#[test]
fn blpop_multi_key_woken_on_second_key() {
let srv = Server::start(1);
let mut consumer = srv.connect();
let mut producer = srv.connect();
consumer
.write_all(&req(&[b"BLPOP", b"k1", b"k2", b"5"]))
.unwrap();
std::thread::sleep(std::time::Duration::from_millis(50));
producer.write_all(&req(&[b"LPUSH", b"k2", b"v2"])).unwrap();
let _ = read_reply(&mut producer); let reply = read_reply(&mut consumer);
assert_eq!(reply, b"*2\r\n$2\r\nk2\r\n$2\r\nv2\r\n");
}
#[test]
fn blpop_multi_key_immediate_hit_first_available() {
let srv = Server::start(1);
let mut c = srv.connect();
c.write_all(&req(&[b"RPUSH", b"ready", b"x"])).unwrap();
let _ = read_reply(&mut c); c.write_all(&req(&[b"BLPOP", b"missing", b"ready", b"5"]))
.unwrap();
let reply = read_reply(&mut c);
assert_eq!(reply, b"*2\r\n$5\r\nready\r\n$1\r\nx\r\n");
}