pub mod distributed;
pub mod lock_manager;
pub mod manager;
pub mod wal;
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use serde_json::Value;
use std::fmt;
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
pub struct TransactionId(u64);
impl TransactionId {
pub fn new() -> Self {
use std::sync::atomic::{AtomicU64, Ordering};
static LAST: AtomicU64 = AtomicU64::new(0);
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos() as u64;
let mut last = LAST.load(Ordering::Relaxed);
loop {
let candidate = now.max(last + 1);
match LAST.compare_exchange_weak(last, candidate, Ordering::Relaxed, Ordering::Relaxed)
{
Ok(_) => return Self(candidate),
Err(actual) => last = actual,
}
}
}
pub fn from_u64(id: u64) -> Self {
Self(id)
}
pub fn as_u64(&self) -> u64 {
self.0
}
}
impl Default for TransactionId {
fn default() -> Self {
Self::new()
}
}
impl fmt::Display for TransactionId {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "tx:{}", self.0)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub enum TransactionState {
Active,
Preparing,
Committed,
Aborted,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Default)]
pub enum IsolationLevel {
ReadUncommitted,
#[default]
ReadCommitted,
RepeatableRead,
Serializable,
}
impl IsolationLevel {
pub fn requires_wal(&self) -> bool {
match self {
IsolationLevel::ReadUncommitted => false,
IsolationLevel::ReadCommitted => true,
IsolationLevel::RepeatableRead => true,
IsolationLevel::Serializable => true,
}
}
}
use crate::driver::protocol::IsolationLevel as ClientIsolationLevel;
impl From<ClientIsolationLevel> for IsolationLevel {
fn from(level: ClientIsolationLevel) -> Self {
match level {
ClientIsolationLevel::ReadCommitted => IsolationLevel::ReadCommitted,
ClientIsolationLevel::RepeatableRead => IsolationLevel::RepeatableRead,
ClientIsolationLevel::Serializable => IsolationLevel::Serializable,
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub enum Operation {
Insert {
database: String,
collection: String,
key: String,
data: Value,
},
Update {
database: String,
collection: String,
key: String,
old_data: Value,
new_data: Value,
},
Delete {
database: String,
collection: String,
key: String,
old_data: Value,
},
PutBlobChunk {
database: String,
collection: String,
key: String,
chunk_index: u32,
data: Vec<u8>,
},
DeleteBlob {
database: String,
collection: String,
key: String,
},
}
impl Operation {
pub fn database(&self) -> &str {
match self {
Operation::Insert { database, .. } => database,
Operation::Update { database, .. } => database,
Operation::Delete { database, .. } => database,
Operation::PutBlobChunk { database, .. } => database,
Operation::DeleteBlob { database, .. } => database,
}
}
pub fn collection(&self) -> &str {
match self {
Operation::Insert { collection, .. } => collection,
Operation::Update { collection, .. } => collection,
Operation::Delete { collection, .. } => collection,
Operation::PutBlobChunk { collection, .. } => collection,
Operation::DeleteBlob { collection, .. } => collection,
}
}
pub fn key(&self) -> &str {
match self {
Operation::Insert { key, .. } => key,
Operation::Update { key, .. } => key,
Operation::Delete { key, .. } => key,
Operation::PutBlobChunk { key, .. } => key,
Operation::DeleteBlob { key, .. } => key,
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Transaction {
pub id: TransactionId,
pub state: TransactionState,
pub isolation_level: IsolationLevel,
pub operations: Vec<Operation>,
pub read_timestamp: u64,
pub write_timestamp: Option<u64>,
pub created_at: DateTime<Utc>,
pub validation_errors: Vec<String>,
}
impl Transaction {
pub fn new(isolation_level: IsolationLevel) -> Self {
let id = TransactionId::new();
let read_timestamp = id.as_u64();
Self {
id,
state: TransactionState::Active,
isolation_level,
operations: Vec::new(),
read_timestamp,
write_timestamp: None,
created_at: Utc::now(),
validation_errors: Vec::new(),
}
}
pub fn add_operation(&mut self, operation: Operation) {
self.operations.push(operation);
}
pub fn is_active(&self) -> bool {
self.state == TransactionState::Active
}
pub fn add_validation_error(&mut self, error: String) {
self.validation_errors.push(error);
}
pub fn has_validation_errors(&self) -> bool {
!self.validation_errors.is_empty()
}
pub fn get_validation_errors(&self) -> &[String] {
&self.validation_errors
}
pub fn clear_validation_errors(&mut self) {
self.validation_errors.clear();
}
pub fn prepare(&mut self) {
self.state = TransactionState::Preparing;
self.write_timestamp = Some(TransactionId::new().as_u64());
}
pub fn commit(&mut self) {
self.state = TransactionState::Committed;
}
pub fn abort(&mut self) {
self.state = TransactionState::Aborted;
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_transaction_id_ordering() {
let id1 = TransactionId::new();
std::thread::sleep(std::time::Duration::from_nanos(100));
let id2 = TransactionId::new();
assert!(id1 < id2);
}
#[test]
fn test_transaction_lifecycle() {
let mut tx = Transaction::new(IsolationLevel::ReadCommitted);
assert_eq!(tx.state, TransactionState::Active);
assert!(tx.is_active());
tx.prepare();
assert_eq!(tx.state, TransactionState::Preparing);
assert!(!tx.is_active());
assert!(tx.write_timestamp.is_some());
tx.commit();
assert_eq!(tx.state, TransactionState::Committed);
}
#[test]
fn test_transaction_operations() {
let mut tx = Transaction::new(IsolationLevel::ReadCommitted);
tx.add_operation(Operation::Insert {
database: "_system".to_string(),
collection: "users".to_string(),
key: "user1".to_string(),
data: serde_json::json!({"name": "Alice"}),
});
assert_eq!(tx.operations.len(), 1);
assert_eq!(tx.operations[0].database(), "_system");
assert_eq!(tx.operations[0].collection(), "users");
assert_eq!(tx.operations[0].key(), "user1");
}
#[test]
fn test_isolation_level_default() {
let level = IsolationLevel::default();
assert_eq!(level, IsolationLevel::ReadCommitted);
}
#[test]
fn test_transaction_id_from_u64() {
let id = TransactionId::from_u64(12345);
assert_eq!(id.as_u64(), 12345);
}
#[test]
fn test_transaction_id_display() {
let id = TransactionId::from_u64(99);
assert_eq!(format!("{}", id), "tx:99");
}
#[test]
fn test_transaction_id_default() {
let id1 = TransactionId::default();
let id2 = TransactionId::default();
assert!(id1.as_u64() > 0);
assert!(id2.as_u64() > 0);
}
#[test]
fn test_transaction_abort() {
let mut tx = Transaction::new(IsolationLevel::Serializable);
tx.abort();
assert_eq!(tx.state, TransactionState::Aborted);
assert!(!tx.is_active());
}
#[test]
fn test_transaction_validation_errors() {
let mut tx = Transaction::new(IsolationLevel::ReadCommitted);
assert!(!tx.has_validation_errors());
tx.add_validation_error("Error 1".to_string());
tx.add_validation_error("Error 2".to_string());
assert!(tx.has_validation_errors());
assert_eq!(tx.get_validation_errors().len(), 2);
tx.clear_validation_errors();
assert!(!tx.has_validation_errors());
}
#[test]
fn test_operation_update() {
let op = Operation::Update {
database: "db".to_string(),
collection: "coll".to_string(),
key: "key1".to_string(),
old_data: serde_json::json!({"a": 1}),
new_data: serde_json::json!({"a": 2}),
};
assert_eq!(op.database(), "db");
assert_eq!(op.collection(), "coll");
assert_eq!(op.key(), "key1");
}
#[test]
fn test_operation_delete() {
let op = Operation::Delete {
database: "mydb".to_string(),
collection: "mycoll".to_string(),
key: "doc1".to_string(),
old_data: serde_json::json!({}),
};
assert_eq!(op.database(), "mydb");
assert_eq!(op.collection(), "mycoll");
assert_eq!(op.key(), "doc1");
}
#[test]
fn test_operation_put_blob_chunk() {
let op = Operation::PutBlobChunk {
database: "blobs".to_string(),
collection: "files".to_string(),
key: "file1".to_string(),
chunk_index: 0,
data: vec![1, 2, 3, 4],
};
assert_eq!(op.database(), "blobs");
assert_eq!(op.collection(), "files");
assert_eq!(op.key(), "file1");
}
#[test]
fn test_operation_delete_blob() {
let op = Operation::DeleteBlob {
database: "blobs".to_string(),
collection: "files".to_string(),
key: "file2".to_string(),
};
assert_eq!(op.database(), "blobs");
assert_eq!(op.key(), "file2");
}
#[test]
fn test_isolation_level_variants() {
let levels = [
IsolationLevel::ReadUncommitted,
IsolationLevel::ReadCommitted,
IsolationLevel::RepeatableRead,
IsolationLevel::Serializable,
];
for level in levels {
let tx = Transaction::new(level);
assert_eq!(tx.isolation_level, level);
}
}
#[test]
fn test_transaction_state_variants() {
let mut tx = Transaction::new(IsolationLevel::ReadCommitted);
assert_eq!(tx.state, TransactionState::Active);
tx.prepare();
assert_eq!(tx.state, TransactionState::Preparing);
tx.commit();
assert_eq!(tx.state, TransactionState::Committed);
}
#[test]
fn test_transaction_serialization() {
let tx = Transaction::new(IsolationLevel::ReadCommitted);
let json = serde_json::to_string(&tx).unwrap();
assert!(json.contains("Active"));
let deserialized: Transaction = serde_json::from_str(&json).unwrap();
assert_eq!(deserialized.id.as_u64(), tx.id.as_u64());
}
}