#![cfg(unix)]
mod common;
use std::time::Duration;
use common::{
encode_connect, encode_query, free_port, read_response, send_sigkill, spawn_server_bin_env,
wait_for_bind, wait_with_timeout,
};
use powdb_server::protocol::Message;
use tokio::io::AsyncWriteExt;
const PRE_KILL_ACKS: usize = 25;
const MAX_INSERTS: usize = 200;
#[tokio::test]
async fn sigkill_mid_write_preserves_every_acknowledged_insert() {
let tmp = tempfile::tempdir().unwrap();
let port = free_port();
let mut child = spawn_server_bin_env(port, tmp.path(), &[], &[("POWDB_SYNC_MODE", "full")]);
let mut stream = wait_for_bind(port, Duration::from_secs(20)).await;
stream.write_all(&encode_connect("testdb")).await.unwrap();
assert_eq!(
read_response(&mut stream).await[0],
0x02,
"expected CONNECT_OK"
);
stream
.write_all(&encode_query("type Ledger { required id: int }"))
.await
.unwrap();
let r = read_response(&mut stream).await;
assert!(r[0] == 0x09 || r[0] == 0x0B, "expected create-type ack");
let mut acked: Vec<i64> = Vec::new();
for i in 0..MAX_INSERTS as i64 {
let frame = encode_query(&format!("insert Ledger {{ id := {i} }}"));
if stream.write_all(&frame).await.is_err() {
break; }
if i == PRE_KILL_ACKS as i64 {
send_sigkill(&child);
}
let resp = tokio::time::timeout(Duration::from_secs(10), async {
let mut header = [0u8; 6];
use tokio::io::AsyncReadExt;
stream.read_exact(&mut header).await?;
let payload_len = u32::from_le_bytes(header[2..6].try_into().unwrap()) as usize;
let mut payload = vec![0u8; payload_len];
if payload_len > 0 {
stream.read_exact(&mut payload).await?;
}
Ok::<u8, std::io::Error>(header[0])
})
.await;
match resp {
Ok(Ok(0x09)) => acked.push(i),
Ok(Ok(other)) => panic!("unexpected reply 0x{other:02X} for insert {i}"),
Ok(Err(_)) | Err(_) => break, }
}
assert!(
acked.len() >= PRE_KILL_ACKS,
"expected at least {PRE_KILL_ACKS} acknowledged inserts before the kill, got {}",
acked.len()
);
let status = wait_with_timeout(&mut child, Duration::from_secs(10));
assert!(!status.success(), "kill -9 must terminate the server");
drop(stream);
let mut child2 = spawn_server_bin_env(port, tmp.path(), &[], &[("POWDB_SYNC_MODE", "full")]);
let mut s2 = wait_for_bind(port, Duration::from_secs(20)).await;
s2.write_all(&encode_connect("testdb")).await.unwrap();
assert_eq!(
read_response(&mut s2).await[0],
0x02,
"expected CONNECT_OK after restart (clean WAL replay)"
);
s2.write_all(&encode_query("Ledger")).await.unwrap();
let resp = read_response(&mut s2).await;
assert_eq!(resp[0], 0x07, "expected RESULT_ROWS after restart");
let recovered: Vec<i64> = match Message::decode(&resp).unwrap() {
Message::ResultRows { columns, rows } => {
let id_col = columns
.iter()
.position(|c| c == "id")
.expect("Ledger must expose the id column");
rows.iter()
.map(|row| row[id_col].parse::<i64>().expect("integer id"))
.collect()
}
other => panic!("expected ResultRows, got {other:?}"),
};
for id in &acked {
assert!(
recovered.contains(id),
"acknowledged insert id={id} lost after kill -9 + restart \
(acked {} rows, recovered {:?})",
acked.len(),
recovered
);
}
s2.write_all(&encode_query("insert Ledger { id := 100000 }"))
.await
.unwrap();
assert_eq!(
read_response(&mut s2).await[0],
0x09,
"restarted server must accept new writes"
);
let _ = child2.kill();
let _ = child2.wait();
}