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);
}