use liven::executor::execute_query;
use liven::storage::StorageEngine;
use liven::types::Query;
use std::sync::Arc;
use tempfile::tempdir;
#[test]
fn test_flusher_shutdown_with_pending_writes() {
let temp_dir = tempdir().unwrap();
let data_dir = temp_dir.path();
let engine = StorageEngine::new(data_dir, 16 * 1024 * 1024).unwrap();
for i in 0..10 {
let _ = engine.append(
"test_stream",
&format!("key_{}", i),
liven::types::DataValue::String(format!("value_{}", i)),
false,
);
}
drop(engine);
let engine2 = StorageEngine::new(data_dir, 16 * 1024 * 1024).unwrap();
let records = engine2.scan_historical().unwrap();
assert_eq!(
records.len(),
10,
"Data should persist after clean shutdown"
);
for i in 0..10 {
let key = format!("key_{}", i);
let record = engine2.get("test_stream", &key).unwrap();
assert!(record.is_some(), "Record {} should exist", i);
}
}
#[test]
fn test_write_lock_prevents_duplicate_keys() {
let temp_dir = tempdir().unwrap();
let data_dir = temp_dir.path();
let engine = Arc::new(StorageEngine::new(data_dir, 16 * 1024 * 1024).unwrap());
let query1 = Query::Insert {
stream_name: "test_stream".to_string(),
key: "test_key".to_string(),
value: serde_json::json!({"field": "value1"}),
};
let result1 = execute_query(&engine, &query1);
assert!(result1.is_ok(), "First insert should succeed");
assert_eq!(result1.unwrap().len(), 1);
let query2 = Query::Insert {
stream_name: "test_stream".to_string(),
key: "test_key".to_string(),
value: serde_json::json!({"field": "value2"}),
};
let result2 = execute_query(&engine, &query2);
assert!(result2.is_err(), "Second insert should fail");
if let Err(e) = result2 {
let error_msg = e.to_string();
assert!(error_msg.contains("already exists"));
assert!(error_msg.contains("upsert"));
assert!(error_msg.contains("update"));
}
let records = engine.get("test_stream", "test_key").unwrap();
assert!(records.is_some());
}
#[test]
fn test_upsert_tombstones_old_records() {
let temp_dir = tempdir().unwrap();
let data_dir = temp_dir.path();
let engine = Arc::new(StorageEngine::new(data_dir, 16 * 1024 * 1024).unwrap());
let insert_query = Query::Insert {
stream_name: "test_stream".to_string(),
key: "test_key".to_string(),
value: serde_json::json!({"version": 1}),
};
execute_query(&engine, &insert_query).unwrap();
let upsert_query = Query::Upsert {
stream_name: "test_stream".to_string(),
key: "test_key".to_string(),
value: serde_json::json!({"version": 2}),
};
execute_query(&engine, &upsert_query).unwrap();
let records = engine.get("test_stream", "test_key").unwrap();
assert!(records.is_some());
let record = records.unwrap();
if let liven::types::DataValue::String(json_str) = &record.value {
let value: serde_json::Value = serde_json::from_str(json_str).unwrap();
assert_eq!(value["version"], 2, "Should have the new version");
} else {
panic!("Expected string JSON value");
}
}
#[test]
fn test_update_tombstones_old_records() {
let temp_dir = tempdir().unwrap();
let data_dir = temp_dir.path();
let engine = Arc::new(StorageEngine::new(data_dir, 16 * 1024 * 1024).unwrap());
let insert_query = Query::Insert {
stream_name: "test_stream".to_string(),
key: "test_key".to_string(),
value: serde_json::json!({"name": "original", "count": 1}),
};
execute_query(&engine, &insert_query).unwrap();
let update_query = Query::Update {
stream_name: "test_stream".to_string(),
key: "test_key".to_string(),
value: serde_json::json!({"count": 2, "new_field": "added"}),
};
execute_query(&engine, &update_query).unwrap();
let records = engine.get("test_stream", "test_key").unwrap();
assert!(records.is_some());
let record = records.unwrap();
if let liven::types::DataValue::String(json_str) = &record.value {
let value: serde_json::Value = serde_json::from_str(json_str).unwrap();
assert_eq!(value["name"], "original", "Should preserve original field");
assert_eq!(value["count"], 2, "Should update count");
assert_eq!(value["new_field"], "added", "Should add new field");
} else {
panic!("Expected string JSON value");
}
}
#[test]
fn test_key_length_validation() {
let temp_dir = tempdir().unwrap();
let data_dir = temp_dir.path();
let engine = Arc::new(StorageEngine::new(data_dir, 16 * 1024 * 1024).unwrap());
let long_key = "a".repeat(33);
let query = Query::Insert {
stream_name: "test_stream".to_string(),
key: long_key,
value: serde_json::json!({"test": "data"}),
};
let result = execute_query(&engine, &query);
assert!(result.is_err(), "33-byte key should be rejected");
if let Err(e) = result {
let error_msg = e.to_string();
assert!(error_msg.contains("33 bytes"));
assert!(error_msg.contains("maximum is 32 bytes"));
assert!(error_msg.contains("Shorten the key"));
}
let valid_key = "a".repeat(32);
let query = Query::Insert {
stream_name: "test_stream".to_string(),
key: valid_key,
value: serde_json::json!({"test": "data"}),
};
let result = execute_query(&engine, &query);
assert!(result.is_ok(), "32-byte key should be accepted");
}