use core::mem::MaybeUninit;
use std::ptr::NonNull;
extern crate alloc;
use alloc::sync::Arc;
use remdb::config::WALConfig;
use remdb::platform::*;
use remdb::transaction::*;
use remdb::types::*;
use remdb::*;
use serial_test::serial;
mod common;
use crate::common::platform::TEST_PLATFORM;
use common::{setup_test_db, setup_test_db_with_memory};
fn create_test_table_def() -> TableDef {
TableDef {
id: 0,
name: "test_table".to_string(),
fields: vec![
FieldDef {
name: "id".to_string(),
data_type: DataType::UInt32,
size: 4,
string_length: None,
offset: 0,
primary_key: true,
not_null: true,
unique: true,
auto_increment: true,
default_value: None,
vector_metadata: None,
json_metadata: None,
},
FieldDef {
name: "value".to_string(),
data_type: DataType::Float32,
size: 4,
string_length: None,
offset: 4,
primary_key: false,
not_null: false,
unique: false,
auto_increment: false,
default_value: None,
vector_metadata: None,
json_metadata: None,
},
],
primary_key: vec![0],
secondary_index: None,
secondary_index_type: IndexType::SortedArray,
record_size: 8,
max_records: 100,
version: 1,
created_at: 0,
updated_at: 0,
}
}
static DEFAULT_ALLOCATOR: config::DefaultMemoryAllocator = config::DefaultMemoryAllocator;
static TEST_DB_CONFIG: std::sync::LazyLock<config::DbConfig> = std::sync::LazyLock::new(|| {
config::DbConfig {
tables: vec![create_test_table_def()],
total_memory: 1024 * 1024, low_power_mode_supported: false,
low_power_max_records: None,
default_max_records: 100000,
memory_allocator: &DEFAULT_ALLOCATOR,
wal_config: WALConfig {
log_path: "./wal",
log_mode: config::LogMode::Sync,
checkpoint_interval_ms: 60000,
log_file_size_limit: 16 * 1024 * 1024,
log_prealloc_size: 1 * 1024 * 1024,
log_segment_size: 16 * 1024 * 1024,
retained_checkpoints: 3,
max_consecutive_invalid: 100,
skip_threshold: 1000,
skip_block_size: 1024 * 1024,
max_skip_attempts: 3,
compression_type: remdb::config::WALCompressionType::None,
compression_level: 3,
},
time_series_defaults: time_series::TimeSeriesConfig::DEFAULT,
#[cfg(feature = "pubsub")]
pubsub_config: None,
#[cfg(feature = "ha")]
ha_config: Some(config::HAConfig {
node_id: 1,
ha_role: remdb::ha::HARole::Auto,
replication_mode: remdb::ha::ReplicationMode::Async,
heartbeat_interval_ms: 1000,
failure_detection_ms: 3000,
sync_timeout_ms: 2000,
master_address: None,
master_port: None,
replication_port: 5556,
}),
model_worker_config: Default::default(),
}
});
static mut TABLES_BUFFER: [Option<MemoryTable>; 1] = [None];
static mut PRIMARY_INDICES_BUFFER: [Option<PrimaryIndex>; 1] = [None];
static mut SECONDARY_INDICES_BUFFER: [Option<AnySecondaryIndex>; 1] = [None];
static mut TABLE_DATA_BUFFER: [u8; 8 * 100] = [0u8; 8 * 100]; static mut TABLE_STATUS_BUFFER: [MaybeUninit<RecordHeader>; 100] =
[const { MaybeUninit::uninit() }; 100];
static mut TABLE_FREE_SLOTS_BUFFER: [usize; 100] = [0usize; 100];
#[test]
#[serial]
fn test_transaction_begin_commit() {
setup_test_db();
remdb::reset_global_db();
crate::transaction::init_tx_manager();
unsafe {
TABLES_BUFFER[0] = None;
PRIMARY_INDICES_BUFFER[0] = None;
SECONDARY_INDICES_BUFFER[0] = None;
TABLE_DATA_BUFFER.fill(0);
for i in 0..100 {
TABLE_STATUS_BUFFER[i].write(RecordHeader {
status: RecordStatus::Free,
version: 0,
lock_type: LockType::None,
lock_owner: 0,
lock_count: 0,
create_tx_id: 0,
delete_tx_id: 0,
next_version_ptr: 0,
});
}
let db = init_global_db(&TEST_DB_CONFIG).unwrap();
let mut tx_buffer = Transaction::default();
let mut log_buffer = Vec::with_capacity(10);
for _ in 0..10 {
log_buffer.push(VariableSizeLogItem::default());
}
let tx = unsafe {
db.begin_transaction(
TransactionType::ReadWrite,
IsolationLevel::ReadCommitted,
&mut tx_buffer,
log_buffer.as_mut_ptr(),
10,
)
}
.unwrap();
let mut record_data = [0u8; 8];
let id: i32 = 1;
let value: f32 = 3.14;
core::ptr::copy_nonoverlapping(&id as *const i32 as *const u8, record_data.as_mut_ptr(), 4);
core::ptr::copy_nonoverlapping(
&value as *const f32 as *const u8,
record_data.as_mut_ptr().add(4),
4,
);
let mut table_mut = db.get_table_mut(0).unwrap();
let record_id = table_mut.insert(record_data.as_ptr()).unwrap();
db.commit_transaction().unwrap();
let table = db.get_table(0).unwrap();
assert_eq!(table.record_count(), 1);
remdb::reset_global_db();
}
}
#[test]
#[serial]
fn test_mvcc_snapshot_isolation() {
setup_test_db();
crate::transaction::init_tx_manager();
unsafe {
TABLES_BUFFER[0] = None;
PRIMARY_INDICES_BUFFER[0] = None;
SECONDARY_INDICES_BUFFER[0] = None;
TABLE_DATA_BUFFER.fill(0);
for i in 0..100 {
TABLE_STATUS_BUFFER[i].write(RecordHeader {
status: RecordStatus::Free,
version: 0,
lock_type: LockType::None,
lock_owner: 0,
lock_count: 0,
create_tx_id: 0,
delete_tx_id: 0,
next_version_ptr: 0,
});
}
let db = init_global_db(&TEST_DB_CONFIG).unwrap();
let mut tx1_buffer = Transaction::default();
let mut tx1_log_buffer = Vec::with_capacity(10);
for _ in 0..10 {
tx1_log_buffer.push(VariableSizeLogItem::default());
}
let tx1 = db
.begin_transaction(
TransactionType::ReadWrite,
IsolationLevel::RepeatableRead,
&mut tx1_buffer,
tx1_log_buffer.as_mut_ptr(),
10,
)
.unwrap();
let mut record_data = [0u8; 8];
let id: u32 = 1;
let value: f32 = 3.14;
core::ptr::copy_nonoverlapping(&id as *const u32 as *const u8, record_data.as_mut_ptr(), 4);
core::ptr::copy_nonoverlapping(
&value as *const f32 as *const u8,
record_data.as_mut_ptr().add(4),
4,
);
let record_id = db
.get_table_mut(0)
.unwrap()
.insert(record_data.as_ptr())
.unwrap();
db.commit_transaction().unwrap();
let mut tx2_buffer = Transaction::default();
let mut tx2_log_buffer = Vec::with_capacity(10);
for _ in 0..10 {
tx2_log_buffer.push(VariableSizeLogItem::default());
}
let tx2 = db
.begin_transaction(
TransactionType::ReadWrite,
IsolationLevel::RepeatableRead,
&mut tx2_buffer,
tx2_log_buffer.as_mut_ptr(),
10,
)
.unwrap();
let mut result_data = [0u8; 8];
{
let table = db.get_table(0).unwrap();
table
.get_by_id(record_id, result_data.as_mut_ptr())
.unwrap();
}
let result_value1 = core::ptr::read(result_data.as_ptr().add(4) as *const f32);
assert_eq!(result_value1, value);
db.commit_transaction().unwrap();
let mut tx3_buffer = Transaction::default();
let mut tx3_log_buffer = Vec::with_capacity(10);
for _ in 0..10 {
tx3_log_buffer.push(VariableSizeLogItem::default());
}
let tx3 = db
.begin_transaction(
TransactionType::ReadWrite,
IsolationLevel::ReadCommitted,
&mut tx3_buffer,
tx3_log_buffer.as_mut_ptr(),
10,
)
.unwrap();
let mut update_data = [0u8; 8];
let new_value: f32 = 6.28;
core::ptr::copy_nonoverlapping(&id as *const u32 as *const u8, update_data.as_mut_ptr(), 4);
core::ptr::copy_nonoverlapping(
&new_value as *const f32 as *const u8,
update_data.as_mut_ptr().add(4),
4,
);
db.get_table_mut(0)
.unwrap()
.update(record_id, update_data.as_ptr())
.unwrap();
db.commit_transaction().unwrap();
let mut tx4_buffer = Transaction::default();
let mut tx4_log_buffer = Vec::with_capacity(10);
for _ in 0..10 {
tx4_log_buffer.push(VariableSizeLogItem::default());
}
let tx4 = db
.begin_transaction(
TransactionType::ReadWrite,
IsolationLevel::ReadCommitted,
&mut tx4_buffer,
tx4_log_buffer.as_mut_ptr(),
10,
)
.unwrap();
{
let table = db.get_table(0).unwrap();
table
.get_by_id(record_id, result_data.as_mut_ptr())
.unwrap();
}
let result_value3 = core::ptr::read(result_data.as_ptr().add(4) as *const f32);
assert_eq!(result_value3, new_value);
db.commit_transaction().unwrap();
remdb::reset_global_db();
}
}
#[test]
#[serial]
fn test_mvcc_version_chain() {
unsafe {
setup_test_db();
crate::transaction::init_tx_manager();
TABLES_BUFFER[0] = None;
PRIMARY_INDICES_BUFFER[0] = None;
SECONDARY_INDICES_BUFFER[0] = None;
TABLE_DATA_BUFFER.fill(0);
for i in 0..100 {
TABLE_STATUS_BUFFER[i].write(RecordHeader {
status: RecordStatus::Free,
version: 0,
lock_type: LockType::None,
lock_owner: 0,
lock_count: 0,
create_tx_id: 0,
delete_tx_id: 0,
next_version_ptr: 0,
});
}
let db = init_global_db(&TEST_DB_CONFIG).unwrap();
let mut record_data = [0u8; 8];
let id: u32 = 1;
let value1: f32 = 1.0;
core::ptr::copy_nonoverlapping(&id as *const u32 as *const u8, record_data.as_mut_ptr(), 4);
core::ptr::copy_nonoverlapping(
&value1 as *const f32 as *const u8,
record_data.as_mut_ptr().add(4),
4,
);
let record_id = db
.get_table_mut(0)
.unwrap()
.insert(record_data.as_ptr())
.unwrap();
let values = [2.0f32, 3.0f32, 4.0f32];
for i in 0..3 {
let mut tx_buffer = Transaction::default();
let mut log_buffer = Vec::with_capacity(10);
for _ in 0..10 {
log_buffer.push(VariableSizeLogItem::default());
}
db.begin_transaction(
TransactionType::ReadWrite,
IsolationLevel::ReadCommitted,
&mut tx_buffer,
log_buffer.as_mut_ptr(),
10,
)
.unwrap();
let mut update_data = [0u8; 8];
let new_value = values[i];
core::ptr::copy_nonoverlapping(
&id as *const u32 as *const u8,
update_data.as_mut_ptr(),
4,
);
core::ptr::copy_nonoverlapping(
&new_value as *const f32 as *const u8,
update_data.as_mut_ptr().add(4),
4,
);
db.get_table_mut(0)
.unwrap()
.update(record_id, update_data.as_ptr())
.unwrap();
db.commit_transaction().unwrap();
}
let mut result_data = [0u8; 8];
let table = db.get_table(0).unwrap();
table
.get_by_id(record_id, result_data.as_mut_ptr())
.unwrap();
let result_value = core::ptr::read(result_data.as_ptr().add(4) as *const f32);
assert_eq!(result_value, 4.0);
let status_ptr = table.status_array.as_ptr().add(record_id);
let current_status = *status_ptr;
assert_eq!(current_status.version, 4);
remdb::reset_global_db();
}
}
#[test]
#[serial]
fn test_mvcc_gc() {
unsafe {
setup_test_db_with_memory(10 * 1024 * 1024);
crate::transaction::init_tx_manager();
TABLES_BUFFER[0] = None;
PRIMARY_INDICES_BUFFER[0] = None;
SECONDARY_INDICES_BUFFER[0] = None;
TABLE_DATA_BUFFER.fill(0);
for i in 0..100 {
TABLE_STATUS_BUFFER[i].write(RecordHeader {
status: RecordStatus::Free,
version: 0,
lock_type: LockType::None,
lock_owner: 0,
lock_count: 0,
create_tx_id: 0,
delete_tx_id: 0,
next_version_ptr: 0,
});
}
let db = init_global_db(&TEST_DB_CONFIG).unwrap();
let mut record_data = [0u8; 8];
let id: u32 = 1;
let value: f32 = 1.0;
core::ptr::copy_nonoverlapping(&id as *const u32 as *const u8, record_data.as_mut_ptr(), 4);
core::ptr::copy_nonoverlapping(
&value as *const f32 as *const u8,
record_data.as_mut_ptr().add(4),
4,
);
let record_id = db
.get_table_mut(0)
.unwrap()
.insert(record_data.as_ptr())
.unwrap();
for i in 0..5 {
let mut tx_buffer = Transaction::default();
let mut log_buffer = Vec::with_capacity(10);
for _ in 0..10 {
log_buffer.push(VariableSizeLogItem::default());
}
db.begin_transaction(
TransactionType::ReadWrite,
IsolationLevel::ReadCommitted,
&mut tx_buffer,
log_buffer.as_mut_ptr(),
10,
)
.unwrap();
let mut update_data = [0u8; 8];
let new_value = (i + 2) as f32;
core::ptr::copy_nonoverlapping(
&id as *const u32 as *const u8,
update_data.as_mut_ptr(),
4,
);
core::ptr::copy_nonoverlapping(
&new_value as *const f32 as *const u8,
update_data.as_mut_ptr().add(4),
4,
);
db.get_table_mut(0)
.unwrap()
.update(record_id, update_data.as_ptr())
.unwrap();
db.commit_transaction().unwrap();
}
let mut result_data = [0u8; 8];
{
let table = db.get_table(0).unwrap();
table
.get_by_id(record_id, result_data.as_mut_ptr())
.unwrap();
}
let result_value = core::ptr::read(result_data.as_ptr().add(4) as *const f32);
assert_eq!(result_value, 6.0);
remdb::reset_global_db();
}
}
#[test]
#[serial]
fn test_mvcc_visibility() {
unsafe {
setup_test_db();
crate::transaction::init_tx_manager();
TABLES_BUFFER[0] = None;
PRIMARY_INDICES_BUFFER[0] = None;
SECONDARY_INDICES_BUFFER[0] = None;
TABLE_DATA_BUFFER.fill(0);
for i in 0..100 {
TABLE_STATUS_BUFFER[i].write(RecordHeader {
status: RecordStatus::Free,
version: 0,
lock_type: LockType::None,
lock_owner: 0,
lock_count: 0,
create_tx_id: 0,
delete_tx_id: 0,
next_version_ptr: 0,
});
}
let db = init_global_db(&TEST_DB_CONFIG).unwrap();
let mut record_data = [0u8; 8];
let id: u32 = 1;
let value: f32 = 1.0;
core::ptr::copy_nonoverlapping(&id as *const u32 as *const u8, record_data.as_mut_ptr(), 4);
core::ptr::copy_nonoverlapping(
&value as *const f32 as *const u8,
record_data.as_mut_ptr().add(4),
4,
);
let record_id = db
.get_table_mut(0)
.unwrap()
.insert(record_data.as_ptr())
.unwrap();
{
let mut tx1_buffer = Transaction::default();
let mut tx1_log_buffer = Vec::with_capacity(10);
for _ in 0..10 {
tx1_log_buffer.push(VariableSizeLogItem::default());
}
let tx1 = db
.begin_transaction(
TransactionType::ReadWrite,
IsolationLevel::ReadCommitted,
&mut tx1_buffer,
tx1_log_buffer.as_mut_ptr(),
10,
)
.unwrap();
let mut update_data = [0u8; 8];
let new_value: f32 = 2.0;
core::ptr::copy_nonoverlapping(
&id as *const u32 as *const u8,
update_data.as_mut_ptr(),
4,
);
core::ptr::copy_nonoverlapping(
&new_value as *const f32 as *const u8,
update_data.as_mut_ptr().add(4),
4,
);
db.get_table_mut(0)
.unwrap()
.update(record_id, update_data.as_ptr())
.unwrap();
db.commit_transaction().unwrap();
}
{
let mut tx2_buffer = Transaction::default();
let mut tx2_log_buffer = Vec::with_capacity(10);
for _ in 0..10 {
tx2_log_buffer.push(VariableSizeLogItem::default());
}
let tx2 = db
.begin_transaction(
TransactionType::ReadWrite,
IsolationLevel::ReadCommitted,
&mut tx2_buffer,
tx2_log_buffer.as_mut_ptr(),
10,
)
.unwrap();
let mut result_data = [0u8; 8];
{
let table = db.get_table(0).unwrap();
table
.get_by_id(record_id, result_data.as_mut_ptr())
.unwrap();
}
let result_value1 = core::ptr::read(result_data.as_ptr().add(4) as *const f32);
assert_eq!(result_value1, 2.0);
db.commit_transaction().unwrap();
}
remdb::reset_global_db();
}
}
#[test]
#[serial]
fn test_transaction_rollback() {
unsafe {
setup_test_db();
crate::transaction::init_tx_manager();
TABLES_BUFFER[0] = None;
PRIMARY_INDICES_BUFFER[0] = None;
SECONDARY_INDICES_BUFFER[0] = None;
TABLE_DATA_BUFFER.fill(0);
for i in 0..100 {
TABLE_STATUS_BUFFER[i].write(RecordHeader {
status: RecordStatus::Free,
version: 0,
lock_type: LockType::None,
lock_owner: 0,
lock_count: 0,
create_tx_id: 0,
delete_tx_id: 0,
next_version_ptr: 0,
});
}
let db = init_global_db(&TEST_DB_CONFIG).unwrap();
let mut tx_buffer = Transaction::default();
let mut log_buffer = Vec::with_capacity(10);
for _ in 0..10 {
log_buffer.push(VariableSizeLogItem::default());
}
let tx = db
.begin_transaction(
TransactionType::ReadWrite,
IsolationLevel::ReadCommitted,
&mut tx_buffer,
log_buffer.as_mut_ptr(),
10,
)
.unwrap();
let mut record_data = [0u8; 8];
let id: i32 = 1;
let value: f32 = 3.14;
core::ptr::copy_nonoverlapping(&id as *const i32 as *const u8, record_data.as_mut_ptr(), 4);
core::ptr::copy_nonoverlapping(
&value as *const f32 as *const u8,
record_data.as_mut_ptr().add(4),
4,
);
let mut table_mut = db.get_table_mut(0).unwrap();
let _record_id = table_mut.insert(record_data.as_ptr()).unwrap();
db.rollback_transaction().unwrap();
let table = db.get_table(0).unwrap();
assert_eq!(table.record_count(), 0);
remdb::reset_global_db();
}
}
#[test]
#[serial]
fn test_transaction_update_rollback() {
unsafe {
setup_test_db();
crate::transaction::init_tx_manager();
TABLES_BUFFER[0] = None;
PRIMARY_INDICES_BUFFER[0] = None;
SECONDARY_INDICES_BUFFER[0] = None;
TABLE_DATA_BUFFER.fill(0);
for i in 0..100 {
TABLE_STATUS_BUFFER[i].write(RecordHeader {
status: RecordStatus::Free,
version: 0,
lock_type: LockType::None,
lock_owner: 0,
lock_count: 0,
create_tx_id: 0,
delete_tx_id: 0,
next_version_ptr: 0,
});
}
let db = init_global_db(&TEST_DB_CONFIG).unwrap();
let mut record_data = [0u8; 8];
let id: i32 = 1;
let value: f32 = 3.14;
core::ptr::copy_nonoverlapping(&id as *const i32 as *const u8, record_data.as_mut_ptr(), 4);
core::ptr::copy_nonoverlapping(
&value as *const f32 as *const u8,
record_data.as_mut_ptr().add(4),
4,
);
let record_id = db
.get_table_mut(0)
.unwrap()
.insert(record_data.as_ptr())
.unwrap();
let mut tx_buffer = Transaction::default();
let mut log_buffer = Vec::with_capacity(10);
for _ in 0..10 {
log_buffer.push(VariableSizeLogItem::default());
}
let tx = db
.begin_transaction(
TransactionType::ReadWrite,
IsolationLevel::ReadCommitted,
&mut tx_buffer,
log_buffer.as_mut_ptr(),
10,
)
.unwrap();
let mut update_data = [0u8; 8];
let new_id: i32 = 1;
let new_value: f32 = 6.28;
core::ptr::copy_nonoverlapping(
&new_id as *const i32 as *const u8,
update_data.as_mut_ptr(),
4,
);
core::ptr::copy_nonoverlapping(
&new_value as *const f32 as *const u8,
update_data.as_mut_ptr().add(4),
4,
);
let mut table_mut = db.get_table_mut(0).unwrap();
table_mut.update(record_id, update_data.as_ptr()).unwrap();
db.rollback_transaction().unwrap();
let table = db.get_table(0).unwrap();
let mut result_data = [0u8; 8];
table
.get_by_id(record_id, result_data.as_mut_ptr())
.unwrap();
let result_id = core::ptr::read(result_data.as_ptr() as *const i32);
let result_value = core::ptr::read(result_data.as_ptr().add(4) as *const f32);
assert_eq!(result_id, id);
assert_eq!(result_value, value);
remdb::reset_global_db();
}
}
#[test]
#[serial]
fn test_transaction_delete_rollback() {
unsafe {
setup_test_db();
crate::transaction::init_tx_manager();
TABLES_BUFFER[0] = None;
PRIMARY_INDICES_BUFFER[0] = None;
SECONDARY_INDICES_BUFFER[0] = None;
TABLE_DATA_BUFFER.fill(0);
for i in 0..100 {
TABLE_STATUS_BUFFER[i].write(RecordHeader {
status: RecordStatus::Free,
version: 0,
lock_type: LockType::None,
lock_owner: 0,
lock_count: 0,
create_tx_id: 0,
delete_tx_id: 0,
next_version_ptr: 0,
});
}
let db = init_global_db(&TEST_DB_CONFIG).unwrap();
let mut record_data = [0u8; 8];
let id: i32 = 1;
let value: f32 = 3.14;
core::ptr::copy_nonoverlapping(&id as *const i32 as *const u8, record_data.as_mut_ptr(), 4);
core::ptr::copy_nonoverlapping(
&value as *const f32 as *const u8,
record_data.as_mut_ptr().add(4),
4,
);
let record_id = db
.get_table_mut(0)
.unwrap()
.insert(record_data.as_ptr())
.unwrap();
let mut tx_buffer = Transaction::default();
let mut log_buffer = Vec::with_capacity(10);
for _ in 0..10 {
log_buffer.push(VariableSizeLogItem::default());
}
let tx = db
.begin_transaction(
TransactionType::ReadWrite,
IsolationLevel::ReadCommitted,
&mut tx_buffer,
log_buffer.as_mut_ptr(),
10,
)
.unwrap();
let mut table_mut = db.get_table_mut(0).unwrap();
table_mut.delete(record_id).unwrap();
db.rollback_transaction().unwrap();
let table = db.get_table(0).unwrap();
assert_eq!(table.record_count(), 1);
let mut result_data = [0u8; 8];
let result = table.get_by_id(record_id, result_data.as_mut_ptr());
assert!(result.is_ok());
remdb::reset_global_db();
}
}