use liven::storage::{StorageEngine, deserialize_payload, serialize_payload_into};
use liven::types::DataValue;
use std::fs;
#[test]
fn test_payload_serialization() {
let stream = "logs";
let key = "user_123";
let val = DataValue::String("login_failed".to_string());
let mut payload = Vec::new();
serialize_payload_into(stream, key, &val, &mut payload);
let (s_out, k_out, v_out) = deserialize_payload(&payload).unwrap();
assert_eq!(s_out, stream);
assert_eq!(k_out, key);
assert_eq!(v_out, val);
}
#[test]
fn test_storage_engine_lifecycle() {
let path = format!(
"./data_liven_test_{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
);
let _ = fs::remove_dir_all(&path).unwrap_err();
let engine = StorageEngine::new(&path, 1024 * 1024).unwrap();
let _rec1 = engine
.append(
"test_stream",
"k1",
DataValue::String("v1".to_string()),
false,
)
.unwrap();
let _rec2 = engine
.append("test_stream", "k2", DataValue::Int(100), false)
.unwrap();
assert_eq!(engine.skipmap.len(), 2);
let get_r1 = engine.get("test_stream", "k1").unwrap().unwrap();
assert_eq!(get_r1.value, DataValue::String("v1".to_string()));
let get_r2 = engine.get("test_stream", "k2").unwrap().unwrap();
assert_eq!(get_r2.value, DataValue::Int(100));
let batch = vec![
(
"test_stream".to_string(),
"k3".to_string(),
DataValue::Bool(true),
),
(
"test_stream".to_string(),
"k4".to_string(),
DataValue::UInt(999),
),
];
engine.append_batch(batch).unwrap();
assert_eq!(engine.skipmap.len(), 4);
drop(engine);
let engine_recovered = StorageEngine::new(&path, 1024 * 1024).unwrap();
assert_eq!(engine_recovered.skipmap.len(), 4);
let recovered_r1 = engine_recovered.get("test_stream", "k1").unwrap().unwrap();
assert_eq!(recovered_r1.value, DataValue::String("v1".to_string()));
let recovered_r4 = engine_recovered.get("test_stream", "k4").unwrap().unwrap();
assert_eq!(recovered_r4.value, DataValue::UInt(999));
engine_recovered
.append("test_stream", "k1", DataValue::Null, true)
.unwrap();
assert_eq!(engine_recovered.skipmap.len(), 3);
assert!(engine_recovered.get("test_stream", "k1").unwrap().is_none());
let streams = engine_recovered.list_streams();
assert_eq!(streams, vec!["test_stream".to_string()]);
let keys = engine_recovered.list_keys("test_stream");
assert!(keys.contains(&"k2".to_string()));
assert!(keys.contains(&"k3".to_string()));
assert!(keys.contains(&"k4".to_string()));
assert!(!keys.contains(&"k1".to_string()));
let _ = fs::remove_dir_all(&path);
}
#[test]
fn test_compaction() {
let path = format!(
"./data_liven_test_compaction_{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
);
let _ = fs::remove_dir_all(&path);
let engine = StorageEngine::new(&path, 120).unwrap();
engine
.append("s1", "k1", DataValue::String("a".to_string()), false)
.unwrap();
engine
.append("s1", "k1", DataValue::String("b".to_string()), false)
.unwrap();
engine
.append("s1", "k2", DataValue::String("c".to_string()), false)
.unwrap();
engine
.append("s1", "k2", DataValue::String("d".to_string()), false)
.unwrap();
engine
.append("s1", "k3", DataValue::String("e".to_string()), false)
.unwrap();
engine
.append("s1", "k3", DataValue::String("f".to_string()), false)
.unwrap();
engine
.append("s1", "k4", DataValue::String("g".to_string()), false)
.unwrap();
engine
.append("s1", "k4", DataValue::String("h".to_string()), false)
.unwrap();
let (_ram_usage, _size_before, segments_before, _streams) = engine.metrics().unwrap();
assert!(segments_before > 1);
engine.compact().unwrap();
let r1 = engine.get("s1", "k1").unwrap().unwrap();
assert_eq!(r1.value, DataValue::String("b".to_string()));
let r2 = engine.get("s1", "k2").unwrap().unwrap();
assert_eq!(r2.value, DataValue::String("d".to_string()));
let r3 = engine.get("s1", "k3").unwrap().unwrap();
assert_eq!(r3.value, DataValue::String("f".to_string()));
let r4 = engine.get("s1", "k4").unwrap().unwrap();
assert_eq!(r4.value, DataValue::String("h".to_string()));
let (_ram_usage_, _size_after, segments_after, _streams_after) = engine.metrics().unwrap();
assert!(segments_after < segments_before);
let _ = fs::remove_dir_all(&path);
}
#[test]
fn test_stream_limits() {
let path = std::env::temp_dir().join(format!(
"liven_test_streams_{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
let _ = fs::remove_dir_all(&path);
let mut engine = StorageEngine::new(&path, 1024 * 1024).unwrap();
engine.set_max_streams(2);
engine.append("s1", "k1", DataValue::Int(1), false).unwrap();
engine.append("s2", "k1", DataValue::Int(1), false).unwrap();
let err = engine.append("s3", "k1", DataValue::Int(1), false);
assert!(err.is_err());
let err_msg = err.unwrap_err().to_string();
assert!(
err_msg.starts_with("stream limit exceeded"),
"got: {}",
err_msg
);
let _ = fs::remove_dir_all(&path);
}
#[test]
fn test_index_ram_limits() {
let path = std::env::temp_dir().join(format!(
"liven_test_ram_{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
let _ = fs::remove_dir_all(&path);
let mut engine = StorageEngine::new(&path, 1024 * 1024).unwrap();
engine.set_max_index_ram_bytes(150);
engine.append("s1", "k1", DataValue::Int(1), false).unwrap();
engine.append("s1", "k2", DataValue::Int(1), false).unwrap();
let err = engine.append("s1", "k3", DataValue::Int(1), false);
assert!(err.is_err());
let err_msg = err.unwrap_err().to_string();
assert!(err_msg.starts_with("Index RAM limit"), "got: {}", err_msg);
engine.append("s1", "k1", DataValue::Null, true).unwrap();
engine.append("s1", "k3", DataValue::Int(1), false).unwrap();
let _ = fs::remove_dir_all(&path);
}
#[test]
fn test_file_descriptor_limits() {
let path = std::env::temp_dir().join(format!(
"liven_test_fds_{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
let _ = fs::remove_dir_all(&path);
let mut engine = StorageEngine::new(&path, 10).unwrap(); engine.set_max_fds(6);
engine
.append("s1", "k1", DataValue::String("a".to_string()), false)
.unwrap(); std::thread::sleep(std::time::Duration::from_millis(150));
engine
.append("s1", "k2", DataValue::String("b".to_string()), false)
.unwrap(); std::thread::sleep(std::time::Duration::from_millis(150));
engine
.append("s1", "k3", DataValue::String("c".to_string()), false)
.unwrap();
let ptr1 = engine
.skipmap
.get(&"s1:k1".to_string())
.map(|e| *e.value())
.unwrap();
let ptr2 = engine
.skipmap
.get(&"s1:k2".to_string())
.map(|e| *e.value())
.unwrap();
let _ = engine.read_record(ptr1).unwrap();
assert!(engine.active_handles_count() <= 1);
let _ = engine.read_record(ptr2).unwrap();
assert!(engine.active_handles_count() <= 1);
let _ = fs::remove_dir_all(&path);
}
#[test]
fn test_scan_excludes_tombstone_frames() {
let path = std::env::temp_dir().join(format!(
"liven_tombstone_scan_{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
let _ = std::fs::remove_dir_all(&path);
let engine = StorageEngine::new(&path, 1024 * 1024).unwrap();
engine
.append("orders", "k1", DataValue::Int(100), false)
.unwrap();
engine
.append("orders", "k2", DataValue::Int(200), false)
.unwrap();
engine
.append("orders", "k3", DataValue::Int(300), false)
.unwrap();
engine
.append("orders", "k2", DataValue::Null, true)
.unwrap();
let records = engine.scan_historical().unwrap();
let order_records: Vec<_> = records
.iter()
.filter(|r| r.stream_name == "orders")
.collect();
assert_eq!(
order_records.len(),
3,
"scan_historical returns active frames (k1, k2, k3); tombstone frame is skipped"
);
assert!(
order_records.iter().all(|r| r.flags & 0x02 == 0),
"No tombstone flags in scan results"
);
assert!(order_records.iter().any(|r| r.key.as_str() == "k1"));
assert!(
order_records.iter().any(|r| r.key.as_str() == "k2"),
"Original k2 active frame must still appear in scan"
);
assert!(order_records.iter().any(|r| r.key.as_str() == "k3"));
let query = liven::parser::parse_query(r#"from("orders") | count()"#).unwrap();
let results = liven::executor::execute_query(&engine, &query).unwrap();
assert_eq!(results.len(), 1);
assert_eq!(
results[0].value,
DataValue::UInt(3),
"count() returns active frames (k1, k2, k3); tombstone is excluded from scan"
);
let _ = std::fs::remove_dir_all(&path);
}
#[test]
fn test_recovery_excludes_fully_tombstoned_streams() {
let path = std::env::temp_dir().join(format!(
"liven_recovery_tombstone_{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
let _ = std::fs::remove_dir_all(&path);
{
let engine = StorageEngine::new(&path, 200).unwrap(); engine
.append("active_stream", "k1", DataValue::Int(1), false)
.unwrap();
engine
.append("dropped_stream", "k1", DataValue::Int(1), false)
.unwrap();
engine
.append("dropped_stream", "k2", DataValue::Int(2), false)
.unwrap();
engine
.append("dropped_stream", "k1", DataValue::Null, true)
.unwrap();
engine
.append("dropped_stream", "k2", DataValue::Null, true)
.unwrap();
engine
.append(
"active_stream",
"padding",
DataValue::String("x".repeat(200)),
false,
)
.unwrap();
engine.compact().unwrap();
}
{
let engine = StorageEngine::new(&path, 1024 * 1024).unwrap();
let streams = engine.list_streams();
assert!(
streams.contains(&"active_stream".to_string()),
"active_stream must survive restart"
);
assert!(
!streams.contains(&"dropped_stream".to_string()),
"dropped_stream has only tombstones — must not appear after restart"
);
assert!(engine.get("active_stream", "k1").unwrap().is_some());
assert!(engine.get("dropped_stream", "k1").unwrap().is_none());
assert!(engine.get("dropped_stream", "k2").unwrap().is_none());
}
let _ = std::fs::remove_dir_all(&path);
}
#[test]
fn test_recovery_stream_limit_excludes_tombstoned_streams() {
let path = std::env::temp_dir().join(format!(
"liven_recovery_limit_{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
let _ = std::fs::remove_dir_all(&path);
{
let engine = StorageEngine::new(&path, 200).unwrap(); engine.append("s1", "k1", DataValue::Int(1), false).unwrap();
engine.append("s2", "k1", DataValue::Int(1), false).unwrap();
engine.append("s3", "k1", DataValue::Int(1), false).unwrap();
engine.append("s2", "k1", DataValue::Null, true).unwrap();
engine.append("s3", "k1", DataValue::Null, true).unwrap();
engine
.append("s1", "padding", DataValue::String("x".repeat(200)), false)
.unwrap();
engine.compact().unwrap();
}
{
let engine = StorageEngine::new(&path, 1024 * 1024).unwrap();
let streams = engine.list_streams();
assert!(!streams.contains(&"s2".to_string()));
assert!(!streams.contains(&"s3".to_string()));
}
let _ = std::fs::remove_dir_all(&path);
}
#[test]
fn test_recovery_torn_write_truncation() {
let path = std::env::temp_dir().join(format!(
"liven_torn_write_{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
let _ = std::fs::remove_dir_all(&path);
{
let engine = StorageEngine::new(&path, 1024 * 1024).unwrap();
engine
.append("test", "k1", DataValue::String("v1".to_string()), false)
.unwrap();
engine
.append("test", "k2", DataValue::Int(42), false)
.unwrap();
engine
.append("test", "k3", DataValue::String("v3".to_string()), false)
.unwrap();
}
let mut segment_files: Vec<_> = std::fs::read_dir(&path)
.unwrap()
.filter_map(|e| e.ok())
.filter(|e| e.path().extension().and_then(|s| s.to_str()) == Some("liven"))
.collect();
segment_files.sort_by_key(|e| e.file_name());
if let Some(last_seg) = segment_files.last() {
let seg_path = last_seg.path();
let original_len = std::fs::metadata(&seg_path).unwrap().len();
assert!(original_len > 100, "Need enough data; got {}", original_len);
let trunc_len = 50u64;
let file = std::fs::OpenOptions::new()
.write(true)
.open(&seg_path)
.unwrap();
file.set_len(trunc_len).unwrap();
}
let engine = StorageEngine::new(&path, 1024 * 1024).unwrap();
assert!(engine.recovery_report.is_some());
let report = engine.recovery_report.as_ref().unwrap();
assert!(
report.had_issues,
"Expected recovery issues after torn write"
);
assert!(
!report.truncated_segments.is_empty(),
"Expected truncation, got: {:?}",
report
);
let k1 = engine.get("test", "k1").unwrap();
assert!(k1.is_some(), "k1 should survive truncation");
let _ = std::fs::remove_dir_all(&path);
}
#[test]
fn test_recovery_mid_file_corruption_skips_segment() {
let path = std::env::temp_dir().join(format!(
"liven_mid_corrupt_{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
let _ = std::fs::remove_dir_all(&path);
{
let engine = StorageEngine::new(&path, 120).unwrap();
engine
.append("stream_a", "a1", DataValue::Int(1), false)
.unwrap();
engine
.append("stream_a", "a2", DataValue::Int(2), false)
.unwrap();
engine
.append("stream_b", "b1", DataValue::Int(3), false)
.unwrap();
engine
.append("stream_b", "b2", DataValue::Int(4), false)
.unwrap();
}
let mut segment_files: Vec<_> = std::fs::read_dir(&path)
.unwrap()
.filter_map(|e| e.ok())
.filter(|e| e.path().extension().and_then(|s| s.to_str()) == Some("liven"))
.collect();
segment_files.sort_by_key(|e| e.file_name());
if segment_files.len() >= 2 {
let first_path = segment_files.first().unwrap().path();
let mut data = std::fs::read(&first_path).unwrap();
if data.len() > 30 {
data[22] ^= 0xFF;
data[23] ^= 0xFF;
std::fs::write(&first_path, &data).unwrap();
}
}
let engine = StorageEngine::new(&path, 1024 * 1024).unwrap();
assert!(engine.recovery_report.is_some());
let report = engine.recovery_report.as_ref().unwrap();
assert!(
report.had_issues,
"Expected recovery issues after corruption"
);
assert!(
!report.skipped_segments.is_empty() || !report.truncated_segments.is_empty(),
"Expected skipped or truncated segments, got: {:?}",
report
);
let _ = std::fs::remove_dir_all(&path);
}
#[test]
fn test_recovery_clean_startup_no_report_issues() {
let path = std::env::temp_dir().join(format!(
"liven_clean_{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
let _ = std::fs::remove_dir_all(&path);
{
let engine = StorageEngine::new(&path, 1024 * 1024).unwrap();
engine
.append("test", "k1", DataValue::Int(100), false)
.unwrap();
}
let engine = StorageEngine::new(&path, 1024 * 1024).unwrap();
assert!(engine.recovery_report.is_some());
let report = engine.recovery_report.as_ref().unwrap();
assert!(!report.had_issues, "Clean startup should report no issues");
assert!(report.truncated_segments.is_empty());
assert!(report.skipped_segments.is_empty());
assert!(report.repair_backups.is_empty());
let _ = std::fs::remove_dir_all(&path);
}