use std::io::{Read, Write};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};
fn dispatch<A: kevy_rt::ArgvView + ?Sized>(store: &mut kevy_store::Store, args: &A) -> Vec<u8> {
thread_local! {
static KEVY: kevy::KevyCommands = kevy::KevyCommands::new();
}
KEVY.with(|k| k.dispatch(store, args))
}
static START_GATE: Mutex<()> = Mutex::new(());
fn free_port_block(width: usize) -> u16 {
use std::sync::atomic::{AtomicU16, Ordering};
use std::sync::LazyLock;
const BAND: u16 = 512;
const LO: u16 = 21_000;
const BANDS: u16 = 80;
static BAND_LO: LazyLock<u16> = LazyLock::new(|| {
let mixed = std::process::id().wrapping_mul(2_654_435_761) >> 16;
LO + (mixed % u32::from(BANDS)) as u16 * BAND
});
static NEXT: LazyLock<AtomicU16> = LazyLock::new(|| AtomicU16::new(*BAND_LO));
'retry: loop {
let span = width as u16 + 1;
let base = NEXT.fetch_add(span, Ordering::Relaxed);
if base.checked_add(span).is_none() || base + span >= *BAND_LO + BAND {
NEXT.store(*BAND_LO, Ordering::Relaxed);
continue;
}
for i in 0..=width as u16 {
if std::net::TcpListener::bind(("127.0.0.1", base + i)).is_err() {
continue 'retry;
}
}
return base;
}
}
fn wait_port(port: u16, what: &str) {
let budget = patience();
let deadline = std::time::Instant::now() + budget;
while std::time::Instant::now() < deadline {
if std::net::TcpStream::connect(("127.0.0.1", port)).is_ok() {
return;
}
std::thread::sleep(std::time::Duration::from_millis(5));
}
panic!("{what} never bound port {port} within {budget:?}");
}
fn patience() -> std::time::Duration {
let mult: f64 = std::env::var("KEVY_TEST_PATIENCE")
.ok()
.and_then(|v| v.parse().ok())
.unwrap_or(1.0);
std::time::Duration::from_secs_f64(60.0 * mult)
}
fn try_connect(port: u16) -> Option<std::net::TcpStream> {
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
while std::time::Instant::now() < deadline {
if let Ok(s) = std::net::TcpStream::connect(("127.0.0.1", port)) {
return Some(s);
}
std::thread::sleep(std::time::Duration::from_millis(5));
}
None
}
struct Server {
#[allow(dead_code)]
port: u16,
replication_base: u16,
nshards: usize,
stop: Arc<AtomicBool>,
handle: Option<std::thread::JoinHandle<()>>,
_dir: Option<TmpDir>,
}
use kevy_tmpdir::TmpDir;
impl Server {
fn start(nshards: usize) -> Server {
Self::start_in(nshards, TmpDir::new("kevy-replication-test"))
}
fn start_in(nshards: usize, dir: TmpDir) -> Server {
let _gate = START_GATE.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
let base = free_port_block(nshards);
let port = base;
let replication_base = base + 1;
let dir_path = dir.path().to_path_buf();
let stop = Arc::new(AtomicBool::new(false));
let stop_thread = stop.clone();
unsafe {
std::env::set_var("KEVY_IO_URING", "0");
}
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_path)
.with_aof(false)
.with_replication(true, 1024 * 1024)
.with_replication_reconnect_window(patience().as_millis() as u32)
.with_replication_listener(replication_base);
let _ = rt.run(stop_thread);
});
let mut ports = vec![port];
ports.extend((0..nshards as u16).map(|i| replication_base + i));
for p in ports {
wait_port(p, "runtime");
}
Server {
port,
replication_base,
nshards,
stop,
handle: Some(handle),
_dir: Some(dir),
}
}
fn shutdown(mut self) {
self.stop.store(true, Ordering::Relaxed);
if let Some(h) = self.handle.take() {
let _ = h.join();
}
}
fn stop_take_dir(mut self) -> TmpDir {
self.stop.store(true, Ordering::Relaxed);
if let Some(h) = self.handle.take() {
let _ = h.join();
}
self._dir.take().expect("dir still owned")
}
}
fn live_generation(replication_port: u16) -> u64 {
let probe = kevy_replicate::replica::ReplicaClient::connect(
("127.0.0.1", replication_port),
"gen-probe",
0,
)
.expect("probe handshake");
probe.primary_gen_at_handshake()
}
fn replicate_from(generation: u64, offset: &str, id: &str) -> Vec<u8> {
let gen_str = generation.to_string();
let mut v = Vec::new();
v.extend_from_slice(b"*6\r\n");
for arg in [
b"REPLICATE".as_slice(),
b"FROM",
gen_str.as_bytes(),
offset.as_bytes(),
b"ID",
id.as_bytes(),
] {
v.extend_from_slice(format!("${}\r\n", arg.len()).as_bytes());
v.extend_from_slice(arg);
v.extend_from_slice(b"\r\n");
}
v
}
fn assert_ack_then_pings(reply: &[u8], want_ack: &[u8]) {
assert!(
reply.starts_with(want_ack),
"reply must start with {:?}, got {:?}",
String::from_utf8_lossy(want_ack),
String::from_utf8_lossy(reply),
);
for line in reply[want_ack.len()..].split(|&b| b == b'\n') {
let line = line.strip_suffix(b"\r").unwrap_or(line);
assert!(
line.is_empty() || line.starts_with(b"+PING "),
"unexpected non-ping bytes after ACK: {:?}",
String::from_utf8_lossy(line),
);
}
}
fn parse_ack(reply: &[u8]) -> (u64, u64, &[u8]) {
let nl = reply
.windows(2)
.position(|w| w == b"\r\n")
.unwrap_or_else(|| panic!("no ACK line in {:?}", String::from_utf8_lossy(reply)));
let line = std::str::from_utf8(&reply[..nl]).expect("ascii ACK");
let mut it = line.split(' ');
assert_eq!(it.next(), Some("+ACK"), "reply must start with an ACK, got {line:?}");
let generation: u64 = it.next().expect("gen field").parse().expect("gen u64");
let offset: u64 = it.next().expect("offset field").parse().expect("offset u64");
(generation, offset, &reply[nl + 2..])
}
fn read_ack(s: &mut std::net::TcpStream) -> (u64, u64, Vec<u8>) {
let mut buf = Vec::new();
while !buf.windows(2).any(|w| w == b"\r\n") {
let mut chunk = [0u8; 256];
match s.read(&mut chunk) {
Ok(0) => break,
Ok(n) => buf.extend_from_slice(&chunk[..n]),
Err(e) => panic!("read while waiting for ACK: {e}"),
}
assert!(buf.len() < 4096, "no ACK line in {:?}", String::from_utf8_lossy(&buf));
}
let (g, o, rest) = parse_ack(&buf);
(g, o, rest.to_vec())
}
fn read_to_eof(s: &mut std::net::TcpStream) -> Vec<u8> {
let _ = s.set_read_timeout(Some(std::time::Duration::from_secs(2)));
let start = std::time::Instant::now();
let mut out = Vec::new();
let mut chunk = [0u8; 256];
while start.elapsed() < std::time::Duration::from_secs(3) {
match s.read(&mut chunk) {
Ok(0) => break,
Ok(n) => out.extend_from_slice(&chunk[..n]),
Err(_) => break,
}
}
out
}
#[test]
fn replica_handshake_receives_ack_and_stays_connected() {
let server = Server::start(1);
let mut s = std::net::TcpStream::connect(("127.0.0.1", server.replication_base)).unwrap();
s.write_all(&replicate_from(live_generation(server.replication_base), "0", "replica-a"))
.unwrap();
let reply = read_to_eof(&mut s);
let (_, off, _) = parse_ack(&reply);
assert_eq!(off, 0);
let ack_len = reply.len() - parse_ack(&reply).2.len();
assert_ack_then_pings(&reply, &reply[..ack_len]);
server.shutdown();
}
#[test]
fn handshake_with_nonzero_offset_echoed_in_ack_then_fence_ships() {
let server = Server::start(1);
let mut s = std::net::TcpStream::connect(("127.0.0.1", server.replication_base)).unwrap();
s.write_all(&replicate_from(0, "12345", "node-7")).unwrap();
let reply = read_to_eof(&mut s);
let (_, off, rest) = parse_ack(&reply);
assert_eq!(off, 12345, "ACK echoes the requested offset");
assert!(
rest.windows(b"+SNAPSHOT\r\n".len()).any(|w| w == b"+SNAPSHOT\r\n"),
"generation fence must ship a snapshot for an unverifiable resume claim, got {:?}",
String::from_utf8_lossy(rest),
);
server.shutdown();
}
#[test]
fn malformed_handshake_closes_connection_no_ack() {
let server = Server::start(1);
let mut s = std::net::TcpStream::connect(("127.0.0.1", server.replication_base)).unwrap();
s.write_all(b"*1\r\n$4\r\nPING\r\n").unwrap();
let reply = read_to_eof(&mut s);
assert!(reply.is_empty(), "got unexpected reply {reply:?}");
server.shutdown();
}
#[test]
fn replication_disabled_means_no_listener_on_replication_port() {
let _gate = START_GATE.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
let base = free_port_block(1);
let dir = TmpDir::new("kevy-replication-disabled");
let dir_path = dir.path().to_path_buf();
let stop = Arc::new(AtomicBool::new(false));
let stop_thread = stop.clone();
unsafe {
std::env::set_var("KEVY_IO_URING", "0");
}
let handle = std::thread::spawn(move || {
let rt = kevy_rt::Runtime::builder(kevy::KevyCommands::sharded(1)).bind([127, 0, 0, 1], base).shards(1)
.with_data_dir(dir_path)
.with_aof(false);
let _ = rt.run(stop_thread);
});
wait_port(base, "server");
let addr: std::net::SocketAddr = format!("127.0.0.1:{}", base + 1).parse().unwrap();
let connect = std::net::TcpStream::connect_timeout(
&addr,
std::time::Duration::from_millis(100),
);
assert!(
connect.is_err(),
"no listener should be on the would-be replication port without with_replication_listener",
);
stop.store(true, Ordering::Relaxed);
let _ = handle.join();
}
fn send_resp(s: &mut std::net::TcpStream, parts: &[&[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");
}
s.write_all(&v).unwrap();
}
fn read_line(s: &mut std::net::TcpStream) -> Vec<u8> {
let _ = s.set_read_timeout(Some(std::time::Duration::from_secs(2)));
let mut line = Vec::new();
let mut b = [0u8; 1];
loop {
s.read_exact(&mut b).unwrap();
line.push(b[0]);
if line.ends_with(b"\r\n") {
return line;
}
}
}
#[test]
fn streaming_replica_receives_set_command_as_wire_frame() {
let server = Server::start(1);
let mut replica = std::net::TcpStream::connect((
"127.0.0.1",
server.replication_base,
))
.unwrap();
replica
.write_all(&replicate_from(live_generation(server.replication_base), "0", "replica-stream"))
.unwrap();
let (_, ack_off, ack_leftover) = read_ack(&mut replica);
assert_eq!(ack_off, 0);
let mut client = std::net::TcpStream::connect(("127.0.0.1", server.port)).unwrap();
send_resp(&mut client, &[b"SET", b"foo", b"bar"]);
let ok = read_line(&mut client);
assert_eq!(ok, b"+OK\r\n");
let mut buf = ack_leftover;
while buf.is_empty() || !buf.windows(2).any(|w| w == b"ar") {
let mut chunk = [0u8; 256];
match replica.read(&mut chunk) {
Ok(0) => break,
Ok(n) => buf.extend_from_slice(&chunk[..n]),
Err(_) => break,
}
if buf.len() > 4096 {
break;
}
}
let mut start = 0usize;
while buf.len() > start && buf[start] == b'+' {
match buf[start..].windows(2).position(|w| w == b"\r\n") {
Some(p) => start += p + 2,
None => break,
}
}
let buf = &buf[start..];
let (offset, argv, used) =
kevy_replicate::wire::decode_frame(buf).expect("decode frame");
assert_eq!(offset, 0);
assert_eq!(argv.len(), 3);
assert_eq!(argv.get(0), Some(&b"SET"[..]));
assert_eq!(argv.get(1), Some(&b"foo"[..]));
assert_eq!(argv.get(2), Some(&b"bar"[..]));
assert!(used <= buf.len());
server.shutdown();
}
#[test]
fn streaming_replica_receives_multiple_frames_in_order() {
let server = Server::start(1);
let mut replica = std::net::TcpStream::connect((
"127.0.0.1",
server.replication_base,
))
.unwrap();
replica
.write_all(&replicate_from(live_generation(server.replication_base), "0", "replica-multi"))
.unwrap();
let (_, ack_off, ack_leftover) = read_ack(&mut replica);
assert_eq!(ack_off, 0);
let mut client = std::net::TcpStream::connect(("127.0.0.1", server.port)).unwrap();
for i in 0..5 {
send_resp(&mut client, &[b"SET", format!("k{i}").as_bytes(), format!("v{i}").as_bytes()]);
let ok = read_line(&mut client);
assert_eq!(ok, b"+OK\r\n");
}
let mut buf = ack_leftover;
let mut frames: Vec<(u64, kevy_resp::Argv)> = Vec::new();
let mut cursor = 0usize;
while frames.len() < 5 {
if buf.len() - cursor > 0 {
if buf[cursor] == b'+'
&& let Some(p) = buf[cursor..].windows(2).position(|w| w == b"\r\n")
{
cursor += p + 2;
continue;
}
match kevy_replicate::wire::decode_frame(&buf[cursor..]) {
Ok((offset, argv, used)) => {
frames.push((offset, argv));
cursor += used;
continue;
}
Err(kevy_replicate::wire::WireError::Truncated) => {
}
Err(e) => panic!("decode error: {e}"),
}
}
let mut chunk = [0u8; 256];
let n = match replica.read(&mut chunk) {
Ok(0) => break,
Ok(n) => n,
Err(_) => break,
};
buf.extend_from_slice(&chunk[..n]);
if buf.len() > 65536 {
break;
}
}
assert_eq!(frames.len(), 5, "expected 5 frames, got {}", frames.len());
for (i, (offset, argv)) in frames.iter().enumerate() {
assert_eq!(*offset, i as u64, "frame {i} offset");
assert_eq!(argv.get(0), Some(&b"SET"[..]));
assert_eq!(argv.get(1), Some(format!("k{i}").as_bytes()));
assert_eq!(argv.get(2), Some(format!("v{i}").as_bytes()));
}
server.shutdown();
}
#[test]
fn streaming_replica_receives_only_its_shards_writes() {
let server = Server::start(2);
let mut replicas: Vec<_> = (0..server.nshards)
.map(|i| {
let port = server.replication_base + i as u16;
let mut r = std::net::TcpStream::connect(("127.0.0.1", port)).unwrap();
r.write_all(&replicate_from(live_generation(port), "0", &format!("replica-{i}")))
.unwrap();
let (_, ack_off, leftover) = read_ack(&mut r);
assert_eq!(ack_off, 0);
(r, leftover)
})
.collect();
let mut client = std::net::TcpStream::connect(("127.0.0.1", server.port)).unwrap();
let keys = ["alpha", "beta", "gamma", "delta", "epsilon", "zeta", "eta", "theta"];
for k in keys {
send_resp(&mut client, &[b"SET", k.as_bytes(), b"v"]);
let ok = read_line(&mut client);
assert_eq!(ok, b"+OK\r\n");
}
let mut total_received = 0usize;
let mut all_keys: Vec<Vec<u8>> = Vec::new();
for (r, leftover) in &mut replicas {
let mut buf = leftover.clone();
let mut cursor = 0usize;
let _ = r.set_read_timeout(Some(std::time::Duration::from_millis(500)));
loop {
if buf.len() > cursor
&& buf[cursor] == b'+'
&& let Some(p) = buf[cursor..].windows(2).position(|w| w == b"\r\n")
{
cursor += p + 2;
continue;
}
match kevy_replicate::wire::decode_frame(&buf[cursor..]) {
Ok((_, argv, used)) => {
cursor += used;
total_received += 1;
all_keys.push(argv.get(1).unwrap().to_vec());
continue;
}
Err(kevy_replicate::wire::WireError::Truncated) => {}
Err(e) => panic!("decode: {e}"),
}
let mut chunk = [0u8; 256];
match r.read(&mut chunk) {
Ok(0) => break,
Ok(n) => buf.extend_from_slice(&chunk[..n]),
Err(_) => break,
}
if buf.len() > 65536 {
break;
}
}
}
assert_eq!(
total_received,
keys.len(),
"expected {} frames across both shards, got {}",
keys.len(),
total_received,
);
all_keys.sort();
let mut expected: Vec<Vec<u8>> = keys.iter().map(|k| k.as_bytes().to_vec()).collect();
expected.sort();
assert_eq!(all_keys, expected);
server.shutdown();
}
#[test]
fn replica_client_handshake_and_receive_set_frame() {
let server = Server::start(1);
let mut client = kevy_replicate::replica::ReplicaClient::connect_at(
("127.0.0.1", server.replication_base),
"replica-via-client",
live_generation(server.replication_base),
0,
std::time::Duration::from_secs(5),
)
.expect("connect + handshake");
assert_eq!(client.primary_offset_at_handshake(), 0);
assert_eq!(client.expected_offset(), 0);
let mut writer = std::net::TcpStream::connect(("127.0.0.1", server.port)).unwrap();
send_resp(&mut writer, &[b"SET", b"foo", b"bar"]);
let ok = read_line(&mut writer);
assert_eq!(ok, b"+OK\r\n");
let frame = client.next().expect("frame").expect("decode ok");
assert_eq!(frame.offset, 0);
assert_eq!(frame.argv.len(), 3);
assert_eq!(frame.argv.get(0), Some(&b"SET"[..]));
assert_eq!(frame.argv.get(1), Some(&b"foo"[..]));
assert_eq!(frame.argv.get(2), Some(&b"bar"[..]));
assert_eq!(client.expected_offset(), 1);
drop(client);
server.shutdown();
}
#[test]
fn replica_client_handshake_failure_on_closed_port() {
let probe = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
let port = probe.local_addr().unwrap().port();
drop(probe);
let result = kevy_replicate::replica::ReplicaClient::connect_with_timeout(
("127.0.0.1", port),
"replica-x",
0,
std::time::Duration::from_millis(200),
);
assert!(
result.is_err(),
"connect to released port should fail, got Ok",
);
}
fn start_small_buffer_primary(buffer_size: u64) -> Server {
let _gate = START_GATE.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
let base = free_port_block(1);
let port = base;
let replication_base = base + 1;
let dir = TmpDir::new("kevy-snapshot-ship");
let dir_path = dir.path().to_path_buf();
let stop = Arc::new(AtomicBool::new(false));
let stop_thread = stop.clone();
unsafe {
std::env::set_var("KEVY_IO_URING", "0");
}
let handle = std::thread::spawn(move || {
let rt = kevy_rt::Runtime::builder(kevy::KevyCommands::sharded(1)).bind([127, 0, 0, 1], port).shards(1)
.with_data_dir(dir_path)
.with_aof(false)
.with_replication(true, buffer_size)
.with_replication_listener(replication_base);
let _ = rt.run(stop_thread);
});
for p in [port, replication_base] {
wait_port(p, "runtime");
}
Server {
port,
replication_base,
nshards: 1,
stop,
handle: Some(handle),
_dir: Some(dir),
}
}
#[test]
fn snapshot_ship_triggers_when_replica_falls_behind_backlog() {
use kevy_replicate::replica::{ReplicaClient, ReplicaEvent};
let server = start_small_buffer_primary(256);
let mut writer = std::net::TcpStream::connect(("127.0.0.1", server.port)).unwrap();
for i in 0..30 {
send_resp(&mut writer, &[b"SET", format!("k{i}").as_bytes(), b"v"]);
let ok = read_line(&mut writer);
assert_eq!(ok, b"+OK\r\n");
}
let mut client = ReplicaClient::connect(
("127.0.0.1", server.replication_base),
"replica-snapshot",
0,
)
.expect("connect + handshake");
loop {
match client.next_event().expect("event").expect("ok") {
ReplicaEvent::SnapshotBegin => break,
kevy_replicate::replica::ReplicaEvent::Ping { .. } => continue,
other => panic!("expected SnapshotBegin, got {other:?}"),
}
}
let mut snapshot_bytes = Vec::new();
let ack_offset = loop {
match client.next_event().expect("event").expect("ok") {
ReplicaEvent::SnapshotChunk(bytes) => snapshot_bytes.extend(bytes),
ReplicaEvent::SnapshotEnd { ack_offset } => break ack_offset,
kevy_replicate::replica::ReplicaEvent::Ping { .. } => continue,
other => panic!("expected SnapshotChunk or SnapshotEnd, got {other:?}"),
}
};
assert_eq!(ack_offset, 30, "ack_offset");
assert_eq!(client.expected_offset(), 30);
assert!(snapshot_bytes.len() > 8, "snapshot too small");
assert_eq!(&snapshot_bytes[..8], b"KEVYSNAP", "snapshot magic");
drop(client);
server.shutdown();
}
#[test]
fn snapshot_ship_loaded_into_local_store_matches_primary() {
use kevy_replicate::replica::{ReplicaClient, ReplicaEvent};
let server = start_small_buffer_primary(256);
let pairs: Vec<(String, String)> = (0..20)
.map(|i| (format!("snap-k{i}"), format!("val-{i:04}")))
.collect();
let mut writer = std::net::TcpStream::connect(("127.0.0.1", server.port)).unwrap();
for (k, v) in &pairs {
send_resp(&mut writer, &[b"SET", k.as_bytes(), v.as_bytes()]);
let ok = read_line(&mut writer);
assert_eq!(ok, b"+OK\r\n");
}
let mut client = ReplicaClient::connect(
("127.0.0.1", server.replication_base),
"replica-loader",
0,
)
.expect("connect");
assert!(matches!(
client.next_event().expect("event").expect("ok"),
ReplicaEvent::SnapshotBegin
));
let mut snapshot_bytes = Vec::new();
let ack_offset = loop {
match client.next_event().expect("event").expect("ok") {
ReplicaEvent::SnapshotChunk(bytes) => snapshot_bytes.extend(bytes),
ReplicaEvent::SnapshotEnd { ack_offset } => break ack_offset,
kevy_replicate::replica::ReplicaEvent::Ping { .. } => continue,
other => panic!("unexpected event: {other:?}"),
}
};
assert_eq!(ack_offset, pairs.len() as u64);
let mut local_store = kevy_store::Store::new();
kevy_persist::load_snapshot_from(&mut local_store, std::io::Cursor::new(&snapshot_bytes))
.expect("load_snapshot_from");
for (k, v) in &pairs {
let argv = kevy::Argv::from(vec![b"GET".to_vec(), k.as_bytes().to_vec()]);
let reply = dispatch(&mut local_store, &argv);
let expected = format!("${}\r\n{}\r\n", v.len(), v);
assert_eq!(
reply, expected.as_bytes(),
"key {k:?}: loaded replica returned {:?}, expected {:?}",
String::from_utf8_lossy(&reply),
expected,
);
}
drop(client);
server.shutdown();
}
#[test]
fn fresh_replica_join_snapshot_then_live_frames() {
use kevy_replicate::replica::{ReplicaClient, ReplicaEvent};
let server = start_small_buffer_primary(256);
let pre: Vec<(String, String)> = (0..20)
.map(|i| (format!("pre-k{i}"), format!("pre-v{i:04}")))
.collect();
let mut writer = std::net::TcpStream::connect(("127.0.0.1", server.port)).unwrap();
for (k, v) in &pre {
send_resp(&mut writer, &[b"SET", k.as_bytes(), v.as_bytes()]);
assert_eq!(read_line(&mut writer), b"+OK\r\n");
}
let mut client = ReplicaClient::connect(
("127.0.0.1", server.replication_base),
"replica-t127",
0,
)
.expect("connect + handshake");
assert!(matches!(
client.next_event().expect("event").expect("ok"),
ReplicaEvent::SnapshotBegin
));
let mut snapshot_bytes = Vec::new();
let ack_offset = loop {
match client.next_event().expect("event").expect("ok") {
ReplicaEvent::SnapshotChunk(bytes) => snapshot_bytes.extend(bytes),
ReplicaEvent::SnapshotEnd { ack_offset } => break ack_offset,
kevy_replicate::replica::ReplicaEvent::Ping { .. } => continue,
other => panic!("expected SnapshotChunk or SnapshotEnd, got {other:?}"),
}
};
assert_eq!(ack_offset, pre.len() as u64);
assert_eq!(client.expected_offset(), ack_offset);
let mut local_store = kevy_store::Store::new();
kevy_persist::load_snapshot_from(&mut local_store, std::io::Cursor::new(&snapshot_bytes))
.expect("load_snapshot_from");
let post: Vec<(String, String)> = (0..5)
.map(|i| (format!("post-k{i}"), format!("post-v{i:04}")))
.collect();
for (k, v) in &post {
send_resp(&mut writer, &[b"SET", k.as_bytes(), v.as_bytes()]);
assert_eq!(read_line(&mut writer), b"+OK\r\n");
}
for (i, _) in post.iter().enumerate() {
let expected_offset = ack_offset + i as u64;
match client.next_event().expect("event").expect("ok") {
ReplicaEvent::Frame(frame) => {
assert_eq!(
frame.offset, expected_offset,
"live frame {i}: offset mismatch (post-snapshot gap)",
);
let _ = dispatch(&mut local_store, &frame.argv);
}
kevy_replicate::replica::ReplicaEvent::Ping { .. } => continue,
other => panic!("live frame {i}: expected Frame, got {other:?}"),
}
}
for (k, v) in pre.iter().chain(post.iter()) {
let argv = kevy::Argv::from(vec![b"GET".to_vec(), k.as_bytes().to_vec()]);
let reply = dispatch(&mut local_store, &argv);
let expected = format!("${}\r\n{}\r\n", v.len(), v);
assert_eq!(
reply, expected.as_bytes(),
"key {k:?}: got {:?}, expected {:?}",
String::from_utf8_lossy(&reply),
expected,
);
}
drop(client);
server.shutdown();
}
#[test]
fn replica_apply_dispatch_mirrors_primary_store() {
let server = Server::start(1);
let mut client = kevy_replicate::replica::ReplicaClient::connect_at(
("127.0.0.1", server.replication_base),
"replica-apply",
live_generation(server.replication_base),
0,
std::time::Duration::from_secs(5),
)
.expect("connect + handshake");
let mut writer = std::net::TcpStream::connect(("127.0.0.1", server.port)).unwrap();
let pairs: &[(&[u8], &[u8])] = &[
(b"alpha", b"one"),
(b"beta", b"two"),
(b"gamma", b"three"),
(b"delta", b"four"),
];
for (k, v) in pairs {
send_resp(&mut writer, &[b"SET", k, v]);
let ok = read_line(&mut writer);
assert_eq!(ok, b"+OK\r\n");
}
let mut local_store = kevy::KeyspaceStore::new();
for expected in 0..pairs.len() as u64 {
let frame = client.next().expect("frame").expect("decode ok");
assert_eq!(frame.offset, expected);
let _reply = dispatch(&mut local_store, &frame.argv);
}
for (k, v) in pairs {
let argv = kevy::Argv::from(vec![b"GET".to_vec(), k.to_vec()]);
let reply = dispatch(&mut local_store, &argv);
let expected = format!("${}\r\n{}\r\n", v.len(), String::from_utf8_lossy(v));
assert_eq!(
reply, expected.as_bytes(),
"key {:?}: replica GET returned {:?}, expected {:?}",
String::from_utf8_lossy(k),
String::from_utf8_lossy(&reply),
expected,
);
}
drop(client);
server.shutdown();
}
#[test]
fn role_reports_master_offset_advancing_with_writes() {
let server = Server::start(1);
let mut s = std::net::TcpStream::connect(("127.0.0.1", server.port)).unwrap();
let mut last = Vec::new();
for _ in 0..500 {
send_resp(&mut s, &[b"ROLE"]);
last = read_line_array(&mut s);
if last.starts_with(b"*3\r\n") {
break;
}
std::thread::sleep(std::time::Duration::from_millis(10));
}
assert!(
last.starts_with(b"*3\r\n$6\r\nmaster\r\n:0\r\n*0\r\n"),
"initial ROLE expected master 0 empty; got {:?}",
String::from_utf8_lossy(&last),
);
for i in 0..7 {
send_resp(&mut s, &[b"SET", format!("rk{i}").as_bytes(), b"v"]);
assert_eq!(read_line(&mut s), b"+OK\r\n");
}
let mut saw_offset = 0u64;
for _ in 0..1000 {
send_resp(&mut s, &[b"ROLE"]);
let reply = read_line_array(&mut s);
if let Some(off) = parse_role_master_offset(&reply) {
saw_offset = off;
if off >= 7 {
break;
}
}
std::thread::sleep(std::time::Duration::from_millis(10));
}
assert_eq!(saw_offset, 7, "ROLE offset should reflect 7 writes");
server.shutdown();
}
fn read_line_array(s: &mut std::net::TcpStream) -> Vec<u8> {
let _ = s.set_read_timeout(Some(std::time::Duration::from_secs(2)));
let mut buf = vec![0u8; 256];
match s.read(&mut buf) {
Ok(n) => buf[..n].to_vec(),
Err(_) => Vec::new(),
}
}
fn parse_role_master_offset(reply: &[u8]) -> Option<u64> {
let prefix = b"*3\r\n$6\r\nmaster\r\n:";
if !reply.starts_with(prefix) {
return None;
}
let rest = &reply[prefix.len()..];
let end = rest.iter().position(|&b| b == b'\r')?;
std::str::from_utf8(&rest[..end]).ok()?.parse().ok()
}
#[test]
fn multi_shard_listener_binds_per_shard_port() {
let server = Server::start(3);
for i in 0..server.nshards {
let port = server.replication_base + i as u16;
let mut s = std::net::TcpStream::connect(("127.0.0.1", port)).unwrap();
s.write_all(&replicate_from(live_generation(port), "0", &format!("replica-{i}")))
.unwrap();
let reply = read_to_eof(&mut s);
let (_, ack_off, rest) = parse_ack(&reply);
assert_eq!(ack_off, 0);
let ack_len = reply.len() - rest.len();
assert_ack_then_pings(&reply, &reply[..ack_len]);
}
server.shutdown();
}
static REPLICA_RT_EXIT: std::sync::Mutex<Option<String>> = std::sync::Mutex::new(None);
struct ReplicaServer {
port: u16,
stop_runtime: Arc<AtomicBool>,
stop_runner: Arc<AtomicBool>,
rt_handle: Option<std::thread::JoinHandle<()>>,
runner_handle: Option<std::thread::JoinHandle<()>>,
_dir: TmpDir,
}
impl ReplicaServer {
fn start(upstream_replication_port: u16) -> Self {
for attempt in 0..5 {
match Self::try_start(upstream_replication_port) {
Ok(s) => return s,
Err(retry_at) => {
eprintln!(
"replica port {retry_at} lost the cross-process bind race \
(attempt {attempt}) — retrying on a fresh port"
);
}
}
}
panic!("replica runtime lost the port race 5 times in a row");
}
fn try_start(upstream_replication_port: u16) -> Result<Self, u16> {
let _gate = START_GATE.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
let port = free_port_block(0) + 1; let dir = TmpDir::new("kevy-replica-rt");
let dir_path = dir.path().to_path_buf();
unsafe {
std::env::set_var("KEVY_IO_URING", "0");
}
let (sender, receiver) = kevy_rt::replica_inbox_pair();
let stop_runtime = Arc::new(AtomicBool::new(false));
let stop_runtime_thread = stop_runtime.clone();
let rt_handle = std::thread::spawn(move || {
let run = std::panic::catch_unwind(std::panic::AssertUnwindSafe(move || {
let rt = kevy_rt::Runtime::builder(kevy::KevyCommands::sharded(1)).bind([127, 0, 0, 1], port).shards(1)
.with_data_dir(dir_path)
.with_aof(false)
.with_replica_inboxes(vec![receiver]);
rt.run(stop_runtime_thread)
}));
let reason = match run {
Ok(Ok(())) => None,
Ok(Err(e)) => Some(format!("Runtime::run returned Err({e})")),
Err(p) => Some(format!(
"Runtime::run PANICKED: {}",
p.downcast_ref::<String>().map(String::as_str).or_else(
|| p.downcast_ref::<&str>().copied()
).unwrap_or("<non-string panic payload>")
)),
};
if let Some(r) = reason {
eprintln!("replica-runtime-exit: {r}");
*REPLICA_RT_EXIT.lock().unwrap_or_else(std::sync::PoisonError::into_inner) =
Some(r);
}
});
{
let deadline = std::time::Instant::now() + patience();
loop {
if std::net::TcpStream::connect(("127.0.0.1", port)).is_ok() {
break;
}
if rt_handle.is_finished() {
let reason = REPLICA_RT_EXIT
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.take();
if reason.as_deref().is_some_and(|r| r.contains("Address already in use")) {
return Err(port);
}
panic!(
"replica runtime exited before binding: {}",
reason.unwrap_or_else(|| "no recorded reason".into())
);
}
if std::time::Instant::now() > deadline {
panic!("replica server never came up on {port}");
}
std::thread::sleep(std::time::Duration::from_millis(10));
}
}
let stop_runner = Arc::new(AtomicBool::new(false));
let stop_runner_thread = stop_runner.clone();
let runner_handle = std::thread::spawn(move || {
let mut from_offset: u64 = 0;
while !stop_runner_thread.load(std::sync::atomic::Ordering::Relaxed) {
let conn = kevy_replicate::replica::ReplicaClient::connect(
("127.0.0.1", upstream_replication_port),
"test-runner",
from_offset,
);
let Ok(mut client) = conn else {
std::thread::sleep(std::time::Duration::from_millis(50));
continue;
};
while !stop_runner_thread.load(std::sync::atomic::Ordering::Relaxed) {
match client.next_event() {
Some(Ok(ev)) => {
let apply = match ev {
kevy_replicate::replica::ReplicaEvent::Ping { .. } => continue,
kevy_replicate::replica::ReplicaEvent::SnapshotBegin => {
kevy_rt::ReplicaApply::SnapshotBegin
}
kevy_replicate::replica::ReplicaEvent::SnapshotChunk(b) => {
kevy_rt::ReplicaApply::SnapshotChunk(b)
}
kevy_replicate::replica::ReplicaEvent::SnapshotEnd { ack_offset } => {
from_offset = ack_offset;
kevy_rt::ReplicaApply::SnapshotEnd { ack_offset, routed: false, gate: None }
}
kevy_replicate::replica::ReplicaEvent::Frame(frame) => {
from_offset = frame.offset.saturating_add(1);
kevy_rt::ReplicaApply::Frame {
offset: frame.offset,
argv: frame.argv,
}
}
};
if sender.send(apply).is_err() {
return;
}
}
Some(Err(_)) | None => break,
}
}
}
});
Ok(Self {
port,
stop_runtime,
stop_runner,
rt_handle: Some(rt_handle),
runner_handle: Some(runner_handle),
_dir: dir,
})
}
fn runtime_alive(&self) -> bool {
self.rt_handle.as_ref().is_some_and(|h| !h.is_finished())
}
fn connect_or_explain(&self, what: &str) -> std::net::TcpStream {
let budget = patience();
let deadline = std::time::Instant::now() + budget;
while std::time::Instant::now() < deadline {
if let Some(s) = try_connect(self.port) {
return s;
}
std::thread::sleep(std::time::Duration::from_millis(5));
}
let state = if self.runtime_alive() {
"the runtime thread is still alive — it bound the port but never \
served in time: a slow runner, widen KEVY_TEST_PATIENCE"
} else {
"the runtime thread has EXITED — the replica panicked rather than \
fell behind; look for its panic, not a timing budget"
};
panic!("{what} never became ready on port {} within {budget:?}. {state}", self.port);
}
fn shutdown(mut self) {
self.stop_runner.store(true, std::sync::atomic::Ordering::Relaxed);
if let Some(h) = self.runner_handle.take() {
let _ = h.join();
}
self.stop_runtime.store(true, std::sync::atomic::Ordering::Relaxed);
if let Some(h) = self.rt_handle.take() {
let _ = h.join();
}
}
}
#[test]
fn server_as_replica_applies_upstream_writes() {
let primary = Server::start(1);
let mut writer = std::net::TcpStream::connect(("127.0.0.1", primary.port)).unwrap();
let pairs: &[(&[u8], &[u8])] = &[
(b"alpha", b"one"),
(b"beta", b"two"),
(b"gamma", b"three"),
(b"delta", b"four"),
(b"epsilon", b"five"),
];
for (k, v) in pairs {
send_resp(&mut writer, &[b"SET", k, v]);
assert_eq!(read_line(&mut writer), b"+OK\r\n");
}
let replica = ReplicaServer::start(primary.replication_base);
let mut reader = replica.connect_or_explain("replica accept loop");
let mut all_seen = false;
for _ in 0..200 {
let mut got_all = true;
for (k, _v) in pairs {
send_resp(&mut reader, &[b"GET", k]);
let line = read_line(&mut reader);
if line.starts_with(b"$-1") {
got_all = false;
break;
}
if line.starts_with(b"$") {
let _ = read_line(&mut reader);
if line.starts_with(b"$0\r") {
eprintln!(
"replica returned $0 (empty bulk) for {:?} mid-catch-up",
String::from_utf8_lossy(k)
);
got_all = false;
break;
}
} else {
got_all = false;
break;
}
}
if got_all {
all_seen = true;
break;
}
std::thread::sleep(std::time::Duration::from_millis(20));
}
assert!(all_seen, "replica did not catch up to all 5 keys in time");
for (k, v) in pairs {
send_resp(&mut reader, &[b"GET", k]);
let header = read_line(&mut reader);
let expected_header = format!("${}\r\n", v.len());
assert_eq!(
header, expected_header.as_bytes(),
"key {:?}: header mismatch", String::from_utf8_lossy(k),
);
let payload = read_line(&mut reader);
let mut expected_payload = v.to_vec();
expected_payload.extend_from_slice(b"\r\n");
assert_eq!(
payload, expected_payload,
"key {:?}: payload mismatch", String::from_utf8_lossy(k),
);
}
drop(reader);
drop(writer);
primary.shutdown();
replica.shutdown();
}
fn read_resp_bulks(s: &mut std::net::TcpStream) -> Vec<Vec<u8>> {
let head = read_line(s);
match head[0] {
b'+' | b':' => Vec::new(),
b'$' => {
let n: i64 = std::str::from_utf8(&head[1..head.len() - 2])
.unwrap()
.parse()
.unwrap();
if n < 0 {
return Vec::new();
}
let mut payload = vec![0u8; n as usize + 2];
s.read_exact(&mut payload).unwrap();
payload.truncate(n as usize);
vec![payload]
}
b'*' => {
let n: i64 = std::str::from_utf8(&head[1..head.len() - 2])
.unwrap()
.parse()
.unwrap();
let mut out = Vec::new();
for _ in 0..n.max(0) {
out.extend(read_resp_bulks(s));
}
out
}
other => panic!(
"unexpected RESP tag {other:?} in {:?}",
String::from_utf8_lossy(&head)
),
}
}
fn smembers_sorted(s: &mut std::net::TcpStream, key: &[u8]) -> Vec<Vec<u8>> {
send_resp(s, &[b"SMEMBERS", key]);
let mut m = read_resp_bulks(s);
m.sort();
m
}
#[test]
fn spop_storm_keeps_replica_sets_identical() {
let primary = Server::start(1);
let mut writer = std::net::TcpStream::connect(("127.0.0.1", primary.port)).unwrap();
let keys: Vec<Vec<u8>> = (0..4).map(|k| format!("spop-set-{k}").into_bytes()).collect();
let all: Vec<Vec<u8>> = (0..50).map(|i| format!("m{i:02}").into_bytes()).collect();
for key in &keys {
let mut argv: Vec<&[u8]> = vec![b"SADD", key];
argv.extend(all.iter().map(Vec::as_slice));
send_resp(&mut writer, &argv);
assert_eq!(read_line(&mut writer), b":50\r\n");
}
let replica = ReplicaServer::start(primary.replication_base);
for _ in 0..10 {
for key in &keys {
send_resp(&mut writer, &[b"SPOP", key]);
let _ = read_resp_bulks(&mut writer);
send_resp(&mut writer, &[b"SPOP", key, b"2"]);
let _ = read_resp_bulks(&mut writer);
}
}
send_resp(&mut writer, &[b"SPOP", &keys[3], b"100"]);
let drained = read_resp_bulks(&mut writer);
assert_eq!(drained.len(), 20, "set 3 should have had 20 members left");
send_resp(&mut writer, &[b"SPOP", &keys[3]]);
assert_eq!(read_line(&mut writer), b"$-1\r\n");
send_resp(&mut writer, &[b"SET", b"spop-fence", b"done"]);
assert_eq!(read_line(&mut writer), b"+OK\r\n");
let mut reader = replica.connect_or_explain("replica accept loop");
let mut fenced = false;
let budget = patience();
let deadline = std::time::Instant::now() + budget;
while std::time::Instant::now() < deadline {
send_resp(&mut reader, &[b"GET", b"spop-fence"]);
if read_resp_bulks(&mut reader) == vec![b"done".to_vec()] {
fenced = true;
break;
}
std::thread::sleep(std::time::Duration::from_millis(20));
}
if !fenced {
let detail = match try_connect(replica.port) {
Some(mut fresh) => {
send_resp(&mut fresh, &[b"DBSIZE"]);
let dbsize = read_line(&mut fresh);
send_resp(&mut fresh, &[b"GET", b"spop-fence"]);
let fence = read_line(&mut fresh);
format!(
"replica DBSIZE = {}\n\
replica GET spop-fence = {}",
String::from_utf8_lossy(&dbsize).trim_end(),
String::from_utf8_lossy(&fence).trim_end(),
)
}
None => format!(
"replica is not accepting connections on port {}; runtime thread {}",
replica.port,
if replica.runtime_alive() {
"still ALIVE (slow, widen KEVY_TEST_PATIENCE)".to_string()
} else {
format!(
"has EXITED — {}",
REPLICA_RT_EXIT
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
.unwrap_or_else(|| "no recorded reason (exit before wrapper?)".to_string())
)
},
),
};
panic!(
"replica never caught up to the post-storm fence within {budget:?}\n{detail}"
);
}
for (i, key) in keys.iter().enumerate() {
let on_primary = smembers_sorted(&mut writer, key);
let on_replica = smembers_sorted(&mut reader, key);
assert_eq!(
on_primary.len(),
if i == 3 { 0 } else { 20 },
"primary set {i}: unexpected survivor count"
);
assert_eq!(
on_replica, on_primary,
"set {i} diverged between primary and replica — SPOP verb replayed?"
);
}
drop(reader);
drop(writer);
primary.shutdown();
replica.shutdown();
}
#[test]
fn replicaof_command_dynamically_attaches_to_primary() {
let primary = Server::start(1);
let mut writer = std::net::TcpStream::connect(("127.0.0.1", primary.port)).unwrap();
let pairs: &[(&[u8], &[u8])] = &[
(b"dy-alpha", b"A"),
(b"dy-beta", b"B"),
(b"dy-gamma", b"C"),
];
for (k, v) in pairs {
send_resp(&mut writer, &[b"SET", k, v]);
assert_eq!(read_line(&mut writer), b"+OK\r\n");
}
let replica_commands = kevy::KevyCommands::sharded(1);
let receivers = replica_commands
.state()
.take_replica_inboxes()
.expect("fresh state");
let replica_port = free_port_block(1) + 1;
let replica_dir = TmpDir::new("kevy-dynamic-replica");
let replica_dir_path = replica_dir.path().to_path_buf();
unsafe { std::env::set_var("KEVY_IO_URING", "0"); }
let replica_stop = Arc::new(AtomicBool::new(false));
let replica_stop_thread = replica_stop.clone();
let replica_handle = std::thread::spawn(move || {
let rt = kevy_rt::Runtime::builder(replica_commands).bind([127, 0, 0, 1], replica_port).shards(1)
.with_data_dir(replica_dir_path)
.with_aof(false)
.with_replica_inboxes(receivers);
let _ = rt.run(replica_stop_thread);
});
wait_port(replica_port, "server");
let mut admin = std::net::TcpStream::connect(("127.0.0.1", replica_port)).unwrap();
send_resp(&mut admin, &[b"ROLE"]);
let role_pre = {
let _ = admin.set_read_timeout(Some(std::time::Duration::from_secs(2)));
let mut buf = vec![0u8; 256];
let n = admin.read(&mut buf).unwrap();
buf[..n].to_vec()
};
assert!(
role_pre.starts_with(b"*3\r\n$6\r\nmaster\r\n"),
"expected master before REPLICAOF; got {:?}",
String::from_utf8_lossy(&role_pre),
);
let primary_port_str = primary.replication_base.to_string();
send_resp(&mut admin, &[b"REPLICAOF", b"127.0.0.1", primary_port_str.as_bytes()]);
let reply = read_line(&mut admin);
assert_eq!(reply, b"+OK\r\n", "REPLICAOF reply: {:?}", String::from_utf8_lossy(&reply));
let mut reader = std::net::TcpStream::connect(("127.0.0.1", replica_port)).unwrap();
let mut all_seen = false;
for _ in 0..200 {
let mut got_all = true;
for (k, _v) in pairs {
send_resp(&mut reader, &[b"GET", k]);
let line = read_line(&mut reader);
if !line.starts_with(b"$") || line.starts_with(b"$-1") {
got_all = false;
break;
}
let _ = read_line(&mut reader);
}
if got_all {
all_seen = true;
break;
}
std::thread::sleep(std::time::Duration::from_millis(20));
}
assert!(all_seen, "replica didn't catch up after dynamic REPLICAOF");
send_resp(&mut admin, &[b"ROLE"]);
let role_during = {
let _ = admin.set_read_timeout(Some(std::time::Duration::from_secs(2)));
let mut buf = vec![0u8; 256];
let n = admin.read(&mut buf).unwrap();
buf[..n].to_vec()
};
assert!(
role_during.starts_with(b"*5\r\n$5\r\nslave\r\n"),
"expected slave after REPLICAOF; got {:?}",
String::from_utf8_lossy(&role_during),
);
send_resp(&mut admin, &[b"REPLICAOF", b"NO", b"ONE"]);
let reply = read_line(&mut admin);
assert_eq!(reply, b"+OK\r\n");
send_resp(&mut admin, &[b"ROLE"]);
let role_after = {
let _ = admin.set_read_timeout(Some(std::time::Duration::from_secs(2)));
let mut buf = vec![0u8; 256];
let n = admin.read(&mut buf).unwrap();
buf[..n].to_vec()
};
assert!(
role_after.starts_with(b"*3\r\n$6\r\nmaster\r\n"),
"expected master after REPLICAOF NO ONE; got {:?}",
String::from_utf8_lossy(&role_after),
);
drop(reader);
drop(admin);
drop(writer);
primary.shutdown();
replica_stop.store(true, std::sync::atomic::Ordering::Relaxed);
let _ = replica_handle.join();
drop(replica_dir);
}
struct AttachedReplica {
port: u16,
stop: Arc<AtomicBool>,
handle: Option<std::thread::JoinHandle<()>>,
_dir: TmpDir,
}
impl AttachedReplica {
fn start(primary_replication_port: u16) -> Self {
let commands = kevy::KevyCommands::sharded(1);
let receivers = commands.state().take_replica_inboxes().expect("fresh state");
let _gate = START_GATE.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
let port = free_port_block(1) + 1;
let dir = TmpDir::new("kevy-v316-replica");
let dir_path = dir.path().to_path_buf();
unsafe { std::env::set_var("KEVY_IO_URING", "0"); }
let stop = Arc::new(AtomicBool::new(false));
let stop_thread = stop.clone();
let handle = std::thread::spawn(move || {
let rt = kevy_rt::Runtime::builder(commands).bind([127, 0, 0, 1], port).shards(1)
.with_data_dir(dir_path)
.with_aof(false)
.with_replica_inboxes(receivers);
let _ = rt.run(stop_thread);
});
wait_port(port, "server");
drop(_gate);
let mut admin = std::net::TcpStream::connect(("127.0.0.1", port)).unwrap();
let port_str = primary_replication_port.to_string();
send_resp(&mut admin, &[b"REPLICAOF", b"127.0.0.1", port_str.as_bytes()]);
assert_eq!(read_line(&mut admin), b"+OK\r\n");
Self { port, stop, handle: Some(handle), _dir: dir }
}
fn shutdown(mut self) {
if let Ok(mut admin) = std::net::TcpStream::connect(("127.0.0.1", self.port)) {
send_resp(&mut admin, &[b"REPLICAOF", b"NO", b"ONE"]);
let _ = read_line(&mut admin);
}
self.stop.store(true, std::sync::atomic::Ordering::Relaxed);
if let Some(h) = self.handle.take() {
let _ = h.join();
}
}
}
fn read_int(s: &mut std::net::TcpStream) -> i64 {
let line = read_line(s);
assert!(line.starts_with(b":"), "expected integer, got {:?}", String::from_utf8_lossy(&line));
std::str::from_utf8(&line[1..line.len() - 2]).unwrap().parse().unwrap()
}
fn try_read_int_array(s: &mut std::net::TcpStream) -> Option<Vec<i64>> {
let header = read_line(s);
if !header.starts_with(b"*") {
return None;
}
let n: usize = std::str::from_utf8(&header[1..header.len() - 2]).ok()?.parse().ok()?;
Some((0..n).map(|_| read_int(s)).collect())
}
fn read_int_array(s: &mut std::net::TcpStream) -> Vec<i64> {
let header = read_line(s);
assert!(header.starts_with(b"*"), "expected array, got {:?}", String::from_utf8_lossy(&header));
let n: usize = std::str::from_utf8(&header[1..header.len() - 2]).unwrap().parse().unwrap();
(0..n).map(|_| read_int(s)).collect()
}
fn wait_replica_gen_learned(replica_port: u16) -> u64 {
let mut c = std::net::TcpStream::connect(("127.0.0.1", replica_port)).unwrap();
for _ in 0..200 {
send_resp(&mut c, &[b"REPL.TOKEN"]);
let Some(pairs) = try_read_int_array(&mut c) else {
std::thread::sleep(std::time::Duration::from_millis(25));
continue;
};
if pairs.len() == 2 && pairs[0] > 0 {
return pairs[0] as u64;
}
std::thread::sleep(std::time::Duration::from_millis(25));
}
panic!("replica never learned the upstream generation from the heartbeat");
}
fn spawn_primary_process(replication_base: u16) -> (kevy_chaos::Harness, u16, std::path::PathBuf) {
let port = kevy_chaos::pick_free_port().expect("primary port");
let dir = std::env::temp_dir().join(format!("kevy-v316-primary-{port}"));
let _ = std::fs::remove_dir_all(&dir);
let cfg = kevy_chaos::HarnessConfig {
kevy_bin: std::path::PathBuf::from(env!("CARGO_BIN_EXE_kevy")),
threads: 1,
..kevy_chaos::HarnessConfig::new(dir.clone(), port)
.with_fsync("everysec")
.with_extra_toml(format!(
"[replication]\nrole = \"primary\"\nlisten_port_base = {replication_base}\n"
))
};
let primary = kevy_chaos::Harness::spawn(cfg).expect("spawn primary kevy");
(primary, port, dir)
}
#[test]
fn wait_with_no_replica_times_out_to_zero_and_wait_zero_is_immediate() {
let server = Server::start(1);
let mut c = std::net::TcpStream::connect(("127.0.0.1", server.port)).unwrap();
send_resp(&mut c, &[b"SET", b"w0", b"v"]);
assert_eq!(read_line(&mut c), b"+OK\r\n");
let t0 = std::time::Instant::now();
send_resp(&mut c, &[b"WAIT", b"1", b"200"]);
assert_eq!(read_int(&mut c), 0);
let elapsed = t0.elapsed();
assert!(
elapsed >= std::time::Duration::from_millis(150),
"WAIT 1 200 with no replica should park ~200ms, returned in {elapsed:?}"
);
let t0 = std::time::Instant::now();
send_resp(&mut c, &[b"WAIT", b"0", b"5000"]);
assert_eq!(read_int(&mut c), 0);
assert!(t0.elapsed() < std::time::Duration::from_secs(1), "WAIT 0 must not park");
server.shutdown();
}
#[test]
fn wait_one_with_live_replica_returns_at_least_one() {
let replication_base = kevy_chaos::pick_free_port().expect("repl port");
let (primary, pport, pdir) = spawn_primary_process(replication_base);
let replica = AttachedReplica::start(replication_base);
let mut c = std::net::TcpStream::connect(("127.0.0.1", pport)).unwrap();
send_resp(&mut c, &[b"SET", b"w1", b"v"]);
assert_eq!(read_line(&mut c), b"+OK\r\n");
let _ = c.set_read_timeout(Some(std::time::Duration::from_secs(8)));
send_resp(&mut c, &[b"WAIT", b"1", b"5000"]);
let n = read_int(&mut c);
assert!(n >= 1, "expected at least 1 acked replica, got {n}");
replica.shutdown();
drop(primary);
let _ = std::fs::remove_dir_all(pdir);
}
#[test]
fn repl_token_on_primary_reports_live_per_shard_pairs() {
let server = Server::start(1);
let mut c = std::net::TcpStream::connect(("127.0.0.1", server.port)).unwrap();
send_resp(&mut c, &[b"REPL.TOKEN"]);
let before = read_int_array(&mut c);
assert_eq!(before.len(), 2, "1 shard → 1 (gen, offset) pair");
assert!(before[0] >= 1, "feed generation starts at ≥1, got {}", before[0]);
send_resp(&mut c, &[b"SET", b"tok", b"v"]);
assert_eq!(read_line(&mut c), b"+OK\r\n");
send_resp(&mut c, &[b"REPL.TOKEN"]);
let after = read_int_array(&mut c);
assert_eq!(after[0], before[0], "generation unchanged by a plain write");
assert!(
after[1] > before[1],
"next_offset must advance past the write: {} → {}",
before[1],
after[1]
);
server.shutdown();
}
#[test]
fn repl_wait_read_your_writes_and_future_token_misdirects() {
let replication_base = kevy_chaos::pick_free_port().expect("repl port");
let (primary, pport, pdir) = spawn_primary_process(replication_base);
let replica = AttachedReplica::start(replication_base);
let _gen = wait_replica_gen_learned(replica.port);
let mut w = std::net::TcpStream::connect(("127.0.0.1", pport)).unwrap();
let mut r = std::net::TcpStream::connect(("127.0.0.1", replica.port)).unwrap();
for round in 0..10 {
let val = format!("v{round}");
send_resp(&mut w, &[b"SET", b"ryw", val.as_bytes()]);
assert_eq!(read_line(&mut w), b"+OK\r\n");
send_resp(&mut w, &[b"REPL.TOKEN"]);
let tok = read_int_array(&mut w);
assert_eq!(tok.len(), 2);
let (g, off) = (tok[0].to_string(), tok[1].to_string());
let _ = r.set_read_timeout(Some(std::time::Duration::from_secs(8)));
send_resp(
&mut r,
&[b"REPL.WAIT", g.as_bytes(), off.as_bytes(), b"TIMEOUT", b"5000"],
);
let reply = read_line(&mut r);
assert_eq!(
reply,
b"+OK\r\n",
"round {round}: REPL.WAIT: {}",
String::from_utf8_lossy(&reply)
);
send_resp(&mut r, &[b"GET", b"ryw"]);
let header = read_line(&mut r);
assert_eq!(header, format!("${}\r\n", val.len()).as_bytes(), "round {round}");
let payload = read_line(&mut r);
assert_eq!(payload, format!("{val}\r\n").as_bytes(), "round {round}");
}
send_resp(&mut w, &[b"REPL.TOKEN"]);
let tok = read_int_array(&mut w);
let (g, off) = (tok[0].to_string(), (tok[1] + 1000).to_string());
let t0 = std::time::Instant::now();
send_resp(
&mut r,
&[b"REPL.WAIT", g.as_bytes(), off.as_bytes(), b"TIMEOUT", b"300"],
);
let reply = read_line(&mut r);
assert!(
reply.starts_with(b"-MISDIRECTED writer is "),
"future token must misdirect, got {}",
String::from_utf8_lossy(&reply)
);
assert!(
t0.elapsed() >= std::time::Duration::from_millis(250),
"future token should park ~TIMEOUT before misdirecting"
);
let bad_gen = (tok[0] + 7).to_string();
let off_now = tok[1].to_string();
let t0 = std::time::Instant::now();
send_resp(
&mut r,
&[b"REPL.WAIT", bad_gen.as_bytes(), off_now.as_bytes(), b"TIMEOUT", b"5000"],
);
let reply = read_line(&mut r);
assert!(
reply.starts_with(b"-MISDIRECTED"),
"gen mismatch must misdirect, got {}",
String::from_utf8_lossy(&reply)
);
assert!(t0.elapsed() < std::time::Duration::from_secs(1), "gen mismatch must not park");
replica.shutdown();
drop(primary);
let _ = std::fs::remove_dir_all(pdir);
}
#[test]
fn unclean_restart_generation_fence_ships_instead_of_aliasing() {
let server = Server::start(1);
let mut client = std::net::TcpStream::connect(("127.0.0.1", server.port)).unwrap();
for i in 0..5 {
send_resp(&mut client, &[b"SET", format!("old{i}").as_bytes(), b"v"]);
assert_eq!(read_line(&mut client), b"+OK\r\n");
}
let probe = kevy_replicate::replica::ReplicaClient::connect(
("127.0.0.1", server.replication_base),
"fence-probe",
0,
)
.expect("probe handshake");
let gen1 = probe.primary_gen_at_handshake();
assert_ne!(gen1, 0, "fresh dir draws a random feed generation");
drop(probe);
drop(client);
let dir = server.stop_take_dir();
let meta = dir.path().join("feed-0.meta");
let _ = std::fs::remove_file(&meta);
let server = Server::start_in(1, dir);
let mut client = std::net::TcpStream::connect(("127.0.0.1", server.port)).unwrap();
for i in 0..10 {
send_resp(&mut client, &[b"SET", format!("new{i}").as_bytes(), b"v"]);
assert_eq!(read_line(&mut client), b"+OK\r\n");
}
let mut replica = kevy_replicate::replica::ReplicaClient::connect_at(
("127.0.0.1", server.replication_base),
"fence-probe",
gen1,
5,
std::time::Duration::from_secs(5),
)
.expect("resume handshake");
assert_ne!(
replica.primary_gen_at_handshake(),
gen1,
"unclean restart must draw a fresh feed generation"
);
loop {
match replica.next_event().expect("event").expect("ok") {
kevy_replicate::replica::ReplicaEvent::SnapshotBegin => break,
kevy_replicate::replica::ReplicaEvent::Ping { .. } => continue,
other => panic!("fence must ship, got {other:?}"),
}
}
let ack_offset = loop {
match replica.next_event().expect("event").expect("ok") {
kevy_replicate::replica::ReplicaEvent::SnapshotChunk(_) => continue,
kevy_replicate::replica::ReplicaEvent::SnapshotEnd { ack_offset } => break ack_offset,
kevy_replicate::replica::ReplicaEvent::Ping { .. } => continue,
other => panic!("unexpected mid-ship event: {other:?}"),
}
};
assert_eq!(ack_offset, 10, "snapshot covers all of the new history");
server.shutdown();
}
#[test]
fn ahead_cursor_ships_snapshot_instead_of_wedging() {
let server = Server::start(1);
let mut client = std::net::TcpStream::connect(("127.0.0.1", server.port)).unwrap();
for i in 0..3 {
send_resp(&mut client, &[b"SET", format!("k{i}").as_bytes(), b"v"]);
assert_eq!(read_line(&mut client), b"+OK\r\n");
}
let probe = kevy_replicate::replica::ReplicaClient::connect(
("127.0.0.1", server.replication_base),
"gen-probe",
0,
)
.expect("probe handshake");
let live_gen = probe.primary_gen_at_handshake();
drop(probe);
let mut replica = kevy_replicate::replica::ReplicaClient::connect_at(
("127.0.0.1", server.replication_base),
"ahead-probe",
live_gen,
999_999,
std::time::Duration::from_secs(5),
)
.expect("ahead handshake");
let mut pings = 0;
loop {
match replica.next_event().expect("event").expect("ok") {
kevy_replicate::replica::ReplicaEvent::SnapshotBegin => break,
kevy_replicate::replica::ReplicaEvent::Ping { .. } => {
pings += 1;
assert!(
pings < 10,
"10 heartbeats and no snapshot — the ahead cursor wedged as caught-up"
);
}
other => panic!("expected a snapshot ship, got {other:?}"),
}
}
server.shutdown();
}
#[test]
fn multi_key_delete_reaches_a_replica() {
let primary = Server::start(1);
let mut w = std::net::TcpStream::connect(("127.0.0.1", primary.port)).unwrap();
for i in 0..6 {
send_resp(&mut w, &[b"SET", format!("k:{i}").as_bytes(), b"v"]);
assert_eq!(read_line(&mut w), b"+OK\r\n");
}
let replica_commands = kevy::KevyCommands::sharded(1);
let receivers = replica_commands.state().take_replica_inboxes().expect("fresh state");
let replica_port = free_port_block(1) + 1;
let replica_dir = TmpDir::new("kevy-replica-multikey-del");
let replica_dir_path = replica_dir.path().to_path_buf();
unsafe { std::env::set_var("KEVY_IO_URING", "0"); }
let replica_stop = Arc::new(AtomicBool::new(false));
let replica_stop_thread = replica_stop.clone();
let replica_handle = std::thread::spawn(move || {
let rt = kevy_rt::Runtime::builder(replica_commands)
.bind([127, 0, 0, 1], replica_port)
.shards(1)
.with_data_dir(replica_dir_path)
.with_aof(false)
.with_replica_inboxes(receivers);
let _ = rt.run(replica_stop_thread);
});
wait_port(replica_port, "server");
let mut admin = std::net::TcpStream::connect(("127.0.0.1", replica_port)).unwrap();
let base = primary.replication_base.to_string();
send_resp(&mut admin, &[b"REPLICAOF", b"127.0.0.1", base.as_bytes()]);
assert_eq!(read_line(&mut admin), b"+OK\r\n");
let mut r = std::net::TcpStream::connect(("127.0.0.1", replica_port)).unwrap();
fn settles(s: &mut std::net::TcpStream, probe: &[&[u8]], want: &[u8]) -> bool {
for _ in 0..250 {
send_resp(s, probe);
if read_line(s) == want {
return true;
}
std::thread::sleep(std::time::Duration::from_millis(20));
}
false
}
assert!(settles(&mut r, &[b"EXISTS", b"k:5"], b":1\r\n"), "replica never caught up");
send_resp(&mut w, &[b"DEL", b"k:1", b"k:2"]);
assert_eq!(read_line(&mut w), b":2\r\n");
assert!(
settles(&mut r, &[b"EXISTS", b"k:1"], b":0\r\n"),
"the replica still holds k:1 after a multi-key DEL on the primary",
);
assert!(settles(&mut r, &[b"EXISTS", b"k:2"], b":0\r\n"), "k:2 too");
assert!(settles(&mut r, &[b"EXISTS", b"k:3"], b":1\r\n"), "k:3 must survive");
replica_stop.store(true, std::sync::atomic::Ordering::SeqCst);
let _ = std::net::TcpStream::connect(("127.0.0.1", replica_port));
let _ = replica_handle.join();
}
#[test]
fn promoted_node_ships_its_keyspace_to_a_fresh_cursor() {
let primary = Server::start(1);
let mut writer = std::net::TcpStream::connect(("127.0.0.1", primary.port)).unwrap();
for i in 0..50 {
send_resp(&mut writer, &[b"SET", format!("pre{i}").as_bytes(), b"v"]);
assert_eq!(read_line(&mut writer), b"+OK\r\n");
}
let commands = kevy::KevyCommands::sharded(1);
let receivers = commands.state().take_replica_inboxes().expect("fresh state");
let base = free_port_block(1);
let node_port = base;
let node_repl_base = base + 1;
let dir = TmpDir::new("kevy-promotion-window");
let dir_path = dir.path().to_path_buf();
unsafe { std::env::set_var("KEVY_IO_URING", "0"); }
let stop = Arc::new(AtomicBool::new(false));
let stop_thread = stop.clone();
let handle = std::thread::spawn(move || {
let rt = kevy_rt::Runtime::builder(commands)
.bind([127, 0, 0, 1], node_port)
.shards(1)
.with_data_dir(dir_path)
.with_aof(false)
.with_replication(true, 1024 * 1024)
.with_replication_listener(node_repl_base)
.with_replica_inboxes(receivers);
let _ = rt.run(stop_thread);
});
wait_port(node_port, "node");
wait_port(node_repl_base, "node replication listener");
let mut admin = std::net::TcpStream::connect(("127.0.0.1", node_port)).unwrap();
let upstream = primary.replication_base.to_string();
send_resp(&mut admin, &[b"REPLICAOF", b"127.0.0.1", upstream.as_bytes()]);
assert_eq!(read_line(&mut admin), b"+OK\r\n");
let mut mirrored = false;
for _ in 0..400 {
send_resp(&mut admin, &[b"GET", b"pre49"]);
let line = read_line(&mut admin);
if line.starts_with(b"$") && !line.starts_with(b"$-1") {
let _ = read_line(&mut admin);
mirrored = true;
break;
}
std::thread::sleep(std::time::Duration::from_millis(20));
}
assert!(mirrored, "the node never mirrored the upstream keyspace");
send_resp(&mut admin, &[b"REPLICAOF", b"NO", b"ONE"]);
assert_eq!(read_line(&mut admin), b"+OK\r\n");
send_resp(&mut admin, &[b"SET", b"postfail", b"v1"]);
assert_eq!(read_line(&mut admin), b"+OK\r\n");
std::thread::sleep(std::time::Duration::from_millis(600));
let mut cursor = kevy_replicate::replica::ReplicaClient::connect(
("127.0.0.1", node_repl_base),
"window-probe",
0,
)
.expect("fresh handshake");
let mut pings = 0;
let outcome = loop {
match cursor.next_event().expect("event").expect("ok") {
kevy_replicate::replica::ReplicaEvent::SnapshotBegin => break "snapshot",
kevy_replicate::replica::ReplicaEvent::Ping { .. } => {
pings += 1;
assert!(
pings < 10,
"10 heartbeats and no snapshot — the fresh cursor sat \
'caught up' at offset 0 while the promoted node's whole \
keyspace (including the promotion-window write) stayed \
invisible: the availgate failover wedge"
);
}
other => panic!("expected a snapshot ship, got {other:?}"),
}
};
assert_eq!(outcome, "snapshot");
stop.store(true, Ordering::SeqCst);
let _ = std::net::TcpStream::connect(("127.0.0.1", node_port));
let _ = handle.join();
primary.shutdown();
}