#![cfg(feature = "std")]
extern crate alloc;
#[cfg(feature = "ha")]
use remdb::config::HAConfig;
use remdb::config::WALConfig;
#[cfg(feature = "ha")]
use remdb::ha::{HARole, ReplicationMode};
use remdb::time_series::TimeSeriesConfig;
use remdb::time_series::TimeSeriesRecord;
use remdb::types::RemDbError;
use remdb::{config, RemDb};
use std::sync::Mutex;
mod common;
use crate::common::platform::TEST_PLATFORM;
use common::{setup_test_db, setup_test_db_with_memory};
static TEST_MUTEX: Mutex<()> = Mutex::new(());
static mut DB_MEMORY: [u8; 104857600] = [0u8; 104857600];
static TEST_DB_CONFIG: std::sync::LazyLock<config::DbConfig> = std::sync::LazyLock::new(|| {
config::DbConfig {
tables: vec![],
total_memory: 104857600,
default_max_records: 100,
low_power_mode_supported: false,
low_power_max_records: None,
memory_allocator: &config::DefaultMemoryAllocator,
wal_config: WALConfig {
log_path: "./wal",
log_mode: config::LogMode::Async,
log_prealloc_size: 0,
log_file_size_limit: 104857600,
log_segment_size: 1048576,
checkpoint_interval_ms: 30000,
retained_checkpoints: 2,
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: TimeSeriesConfig::DEFAULT,
#[cfg(feature = "pubsub")]
pubsub_config: None,
#[cfg(feature = "ha")]
ha_config: Some(HAConfig {
node_id: 1, ha_role: HARole::Auto,
replication_mode: ReplicationMode::Async,
heartbeat_interval_ms: 1000,
failure_detection_ms: 3000,
sync_timeout_ms: 1000,
master_address: None,
master_port: None,
replication_port: 5556,
}),
model_worker_config: Default::default(),
}
});
static PERFORMANCE_TEST_DB_CONFIG: std::sync::LazyLock<config::DbConfig> =
std::sync::LazyLock::new(|| {
config::DbConfig {
tables: vec![],
total_memory: 104857600,
default_max_records: 100000,
low_power_mode_supported: false,
low_power_max_records: None,
memory_allocator: &config::DefaultMemoryAllocator,
wal_config: WALConfig {
log_path: "./wal",
log_mode: config::LogMode::Async,
log_prealloc_size: 0,
log_file_size_limit: 104857600,
log_segment_size: 1048576,
checkpoint_interval_ms: 30000,
retained_checkpoints: 2,
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: TimeSeriesConfig::DEFAULT,
#[cfg(feature = "pubsub")]
pubsub_config: None,
#[cfg(feature = "ha")]
ha_config: Some(HAConfig {
node_id: 1, ha_role: HARole::Auto,
replication_mode: ReplicationMode::Async,
heartbeat_interval_ms: 1000,
failure_detection_ms: 3000,
sync_timeout_ms: 1000,
master_address: None,
master_port: None,
replication_port: 5556,
}),
model_worker_config: Default::default(),
}
});
static ROLLBACK_TEST_DB_CONFIG: std::sync::LazyLock<config::DbConfig> =
std::sync::LazyLock::new(|| {
config::DbConfig {
tables: vec![],
total_memory: 104857600,
default_max_records: 10000,
low_power_mode_supported: false,
low_power_max_records: None,
memory_allocator: &config::DefaultMemoryAllocator,
wal_config: WALConfig {
log_path: "./wal",
log_mode: config::LogMode::Sync,
log_prealloc_size: 0,
log_file_size_limit: 104857600,
log_segment_size: 1048576,
checkpoint_interval_ms: 30000,
retained_checkpoints: 2,
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: TimeSeriesConfig::DEFAULT,
#[cfg(feature = "pubsub")]
pubsub_config: None,
#[cfg(feature = "ha")]
ha_config: Some(HAConfig {
node_id: 1, ha_role: HARole::Auto,
replication_mode: ReplicationMode::Async,
heartbeat_interval_ms: 1000,
failure_detection_ms: 3000,
sync_timeout_ms: 1000,
master_address: None,
master_port: None,
replication_port: 5556,
}),
model_worker_config: Default::default(),
}
});
#[test]
fn test_write_timeseries_batch_acid() {
let _guard = TEST_MUTEX.lock().unwrap();
setup_test_db_with_memory(10 * 1024 * 1024);
remdb::reset_global_db();
let mut db = RemDb::new(&*TEST_DB_CONFIG);
db.init().unwrap();
let table_name = "test_timeseries";
let time_field = "timestamp";
let value_field = "value";
let tag_fields = &["tag1", "tag2"];
db.create_time_series_table(table_name, time_field, value_field, tag_fields, None)
.unwrap();
let mut data_points = Vec::new();
for i in 0..10 {
data_points.push(TimeSeriesRecord {
timestamp: 1000000 + i as u64,
value: i as f64,
tag_count: 2,
tags: [i as u64, (i * 2) as u64, 0, 0, 0, 0, 0, 0],
});
}
let result = db.write_timeseries_batch(table_name, &data_points);
assert!(result.is_ok());
assert_eq!(result.unwrap(), 10);
let query_result = db
.get_time_series_table(0)
.unwrap()
.query_time_range(1000000, 1000009)
.unwrap();
assert_eq!(query_result.len(), 10);
let result = db.write_timeseries_batch(table_name, &[]);
assert!(result.is_err());
assert_eq!(result.err().unwrap(), RemDbError::ConfigError);
let mut large_data_points = Vec::new();
for i in 0..1000 {
large_data_points.push(TimeSeriesRecord {
timestamp: 2000000 + i as u64,
value: i as f64,
tag_count: 2,
tags: [i as u64, (i * 2) as u64, 0, 0, 0, 0, 0, 0],
});
}
let start_time = std::time::Instant::now();
let result = db.write_timeseries_batch(table_name, &large_data_points);
let duration = start_time.elapsed();
assert!(result.is_ok());
assert_eq!(result.unwrap(), 1000);
println!("写入1000个数据点耗时: {:?}", duration);
assert!(
duration.as_millis() < 8,
"写入1000个数据点耗时超过8毫秒: {:?}",
duration
);
let query_result = db
.get_time_series_table(0)
.unwrap()
.query_time_range(2000000, 2000999)
.unwrap();
assert_eq!(query_result.len(), 1000);
println!("测试通过: 事务化批量写入ACID特性验证成功");
}
#[test]
fn test_write_timeseries_batch_performance() {
let _guard = TEST_MUTEX.lock().unwrap();
setup_test_db_with_memory(20 * 1024 * 1024);
remdb::reset_global_db();
let mut db = RemDb::new(&*PERFORMANCE_TEST_DB_CONFIG);
db.init().unwrap();
let table_name = "performance_timeseries";
let time_field = "timestamp";
let value_field = "value";
let tag_fields = &["tag1"];
db.create_time_series_table(table_name, time_field, value_field, tag_fields, None)
.unwrap();
let total_points: usize = 120000;
let batch_size = 1000;
let mut all_data_points = Vec::new();
for i in 0..total_points {
all_data_points.push(TimeSeriesRecord {
timestamp: 3000000 + i as u64,
value: i as f64,
tag_count: 1,
tags: [i as u64, 0, 0, 0, 0, 0, 0, 0],
});
}
let start_time = std::time::Instant::now();
let mut written = 0;
for chunk in all_data_points.chunks(batch_size) {
let result = db.write_timeseries_batch(table_name, chunk);
written += result.unwrap();
}
let duration = start_time.elapsed();
let throughput = written as f64 / duration.as_secs_f64();
println!("写入 {} 个数据点耗时: {:?}", written, duration);
println!("吞吐量: {:.2} 点/秒", throughput);
assert!(
throughput >= 120000.0,
"吞吐量未达到目标: {:.2} 点/秒 < 120000 点/秒",
throughput
);
let query_result = db
.get_time_series_table(0)
.unwrap()
.query_time_range(3000000, 3000000 + (total_points - 1) as u64)
.unwrap();
assert_eq!(query_result.len(), total_points);
println!("性能测试通过: 吞吐量达到 {:.2} 点/秒", throughput);
}
#[test]
fn test_write_timeseries_batch_rollback() {
let _guard = TEST_MUTEX.lock().unwrap();
setup_test_db_with_memory(10 * 1024 * 1024);
remdb::reset_global_db();
let mut db = RemDb::new(&*ROLLBACK_TEST_DB_CONFIG);
db.init().unwrap();
let table_name = "rollback_timeseries";
let time_field = "timestamp";
let value_field = "value";
let tag_fields = &["tag1"];
db.create_time_series_table(table_name, time_field, value_field, tag_fields, None)
.unwrap();
let mut data_points = Vec::new();
for i in 0..10 {
data_points.push(TimeSeriesRecord {
timestamp: 4000000 + i as u64,
value: i as f64,
tag_count: 1,
tags: [i as u64, 0, 0, 0, 0, 0, 0, 0],
});
}
let result = db.write_timeseries_batch(table_name, &data_points);
assert!(result.is_ok());
assert_eq!(result.unwrap(), 10);
let query_result = db
.get_time_series_table(0)
.unwrap()
.query_time_range(4000000, 4000009)
.unwrap();
assert_eq!(query_result.len(), 10, "数据写入失败");
println!("测试通过: 批量写入验证成功");
}
#[test]
fn test_time_type_support() {
let _guard = TEST_MUTEX.lock().unwrap();
setup_test_db_with_memory(10 * 1024 * 1024);
remdb::reset_global_db();
let mut db = RemDb::new(&TEST_DB_CONFIG);
db.init().unwrap();
let create_table_sql = "CREATE TABLE test_time_types (
id INTEGER PRIMARY KEY AUTO_INCREMENT,
ts TIMESTAMP(3),
tstz TIMESTAMPTZ(6),
name TEXT
)";
let result = db.sql_query(create_table_sql);
if !result.is_ok() {
println!(
"CREATE TABLE failed with error: {:?}",
result.as_ref().err()
);
}
assert!(result.is_ok());
let insert_sql = "INSERT INTO test_time_types (ts, tstz, name) VALUES
(NOW(), CURRENT_TIMESTAMP(), 'test1'),
(LOCALTIMESTAMP(), NOW(), 'test2')";
let result = db.sql_query(insert_sql);
if !result.is_ok() {
println!("INSERT failed with error: {:?}", result.as_ref().err());
}
assert!(result.is_ok());
let select_sql = "SELECT id, ts, tstz, name FROM test_time_types";
let result = db.sql_query(select_sql);
assert!(result.is_ok());
let format_sql = "SELECT
TO_ISO8601(ts) as iso_ts,
TO_CHAR(tstz, 'YYYY-MM-DD HH24:MI:SS') as char_tstz,
TO_EPOCH(ts) as epoch_ts
FROM test_time_types";
let result = db.sql_query(format_sql);
if !result.is_ok() {
println!("Format SQL failed with error: {:?}", result.as_ref().err());
}
assert!(result.is_ok());
println!("测试通过: 时间类型和时间格式化函数支持验证成功");
}
#[test]
fn test_time_arithmetic() {
let _guard = TEST_MUTEX.lock().unwrap();
setup_test_db_with_memory(10 * 1024 * 1024);
remdb::reset_global_db();
let mut db = RemDb::new(&TEST_DB_CONFIG);
db.init().unwrap();
let create_table_sql = "CREATE TABLE test_time_arithmetic (
id INTEGER PRIMARY KEY AUTO_INCREMENT,
start_time TIMESTAMP(6),
end_time TIMESTAMPTZ(6)
)";
let result = db.sql_query(create_table_sql);
assert!(result.is_ok());
let insert_sql = "INSERT INTO test_time_arithmetic (start_time, end_time) VALUES
(NOW(), CURRENT_TIMESTAMP())";
let result = db.sql_query(insert_sql);
assert!(result.is_ok());
let add_sql = "SELECT
start_time + INTERVAL 1 HOUR as one_hour_later,
end_time + INTERVAL 30 MINUTE as thirty_minutes_later
FROM test_time_arithmetic";
println!("Executing SQL: {}", add_sql);
let result = db.sql_query(add_sql);
if let Err(e) = &result {
println!("Error executing SQL: {:?}", e);
}
assert!(result.is_ok(), "Error executing SQL: {:?}", result.err());
let sub_sql = "SELECT
start_time - INTERVAL 1 DAY as one_day_ago,
end_time - INTERVAL 1 WEEK as one_week_ago
FROM test_time_arithmetic";
println!("Executing SQL: {}", sub_sql);
let result = db.sql_query(sub_sql);
if let Err(e) = &result {
println!("Error executing SQL: {:?}", e);
}
assert!(result.is_ok());
let diff_sql = "SELECT
end_time - start_time as time_diff
FROM test_time_arithmetic";
println!("Executing SQL: {}", diff_sql);
let result = db.sql_query(diff_sql);
if let Err(e) = &result {
println!("Error executing SQL: {:?}", e);
}
assert!(result.is_ok());
let compare_sql = "SELECT
start_time < end_time as is_start_earlier,
start_time = end_time as is_same,
start_time > end_time as is_start_later
FROM test_time_arithmetic";
println!("Executing SQL: {}", compare_sql);
let result = db.sql_query(compare_sql);
if let Err(e) = &result {
println!("Error executing SQL: {:?}", e);
}
assert!(result.is_ok());
println!("测试通过: 时间运算和比较功能验证成功");
}
#[test]
fn test_time_precision_support() {
let _guard = TEST_MUTEX.lock().unwrap();
setup_test_db_with_memory(10 * 1024 * 1024);
remdb::reset_global_db();
let mut db = RemDb::new(&TEST_DB_CONFIG);
db.init().unwrap();
let create_table_sql = "CREATE TABLE test_time_precision (
id INTEGER PRIMARY KEY AUTO_INCREMENT,
ts_sec TIMESTAMP(0),
ts_ms TIMESTAMP(3),
ts_us TIMESTAMP(6),
ts_ns TIMESTAMP(9),
tstz_sec TIMESTAMPTZ(0),
tstz_ms TIMESTAMPTZ(3),
tstz_us TIMESTAMPTZ(6),
tstz_ns TIMESTAMPTZ(9)
)";
let result = db.sql_query(create_table_sql);
assert!(result.is_ok());
let insert_sql = "INSERT INTO test_time_precision (ts_sec, ts_ms, ts_us, ts_ns, tstz_sec, tstz_ms, tstz_us, tstz_ns) VALUES
(NOW(), NOW(), NOW(), NOW(), NOW(), NOW(), NOW(), NOW())";
let result = db.sql_query(insert_sql);
assert!(result.is_ok());
let select_sql = "SELECT
ts_sec, ts_ms, ts_us, ts_ns,
tstz_sec, tstz_ms, tstz_us, tstz_ns
FROM test_time_precision";
let result = db.sql_query(select_sql);
assert!(result.is_ok());
println!("测试通过: 时间精度支持验证成功");
}
#[test]
fn test_time_series_pre_aggregation() {
let _guard = TEST_MUTEX.lock().unwrap();
setup_test_db_with_memory(10 * 1024 * 1024);
remdb::reset_global_db();
let mut db = RemDb::new(&*TEST_DB_CONFIG);
db.init().unwrap();
let table_name = "test_pre_aggregation";
let time_field = "timestamp";
let value_field = "value";
let tag_fields = &["tag1"];
db.create_time_series_table(table_name, time_field, value_field, tag_fields, None)
.unwrap();
{
let time_series_table = db.get_time_series_table(0).unwrap();
let interval_seconds = 60; let aggregation = "sum";
let result = time_series_table.add_pre_aggregation(interval_seconds, aggregation);
assert!(result.is_ok());
}
let mut data_points = Vec::new();
let base_timestamp = 1000000000000; for i in 0..10 {
data_points.push(TimeSeriesRecord {
timestamp: base_timestamp + (i * 100000000) as u64, value: i as f64,
tag_count: 1,
tags: [1 as u64, 0, 0, 0, 0, 0, 0, 0],
});
}
let time_series_table = db.get_time_series_table_mut(0).unwrap();
unsafe {
let result = time_series_table.batch_write(data_points.as_ptr(), data_points.len());
assert!(result.is_ok());
assert_eq!(result.unwrap(), 10);
}
{
let time_series_table = db.get_time_series_table(0).unwrap();
let start_time = 1000000000000; let end_time = 1000000000000 + 1000000000; let interval_seconds = 60;
let aggregation = "sum";
let result = time_series_table.query_pre_aggregated(
start_time,
end_time,
interval_seconds,
aggregation,
);
assert!(result.is_ok());
let records = result.unwrap();
assert!(!records.is_empty());
println!("预聚合查询结果: {:?}", records);
}
{
let time_series_table = db.get_time_series_table(0).unwrap();
let interval_seconds = 60;
let avg_aggregation = "avg";
let result = time_series_table.add_pre_aggregation(interval_seconds, avg_aggregation);
assert!(result.is_ok());
}
let mut more_data_points = Vec::new();
let base_timestamp = 1000000000000; for i in 10..20 {
more_data_points.push(TimeSeriesRecord {
timestamp: base_timestamp + (i * 100000000) as u64, value: i as f64,
tag_count: 1,
tags: [1 as u64, 0, 0, 0, 0, 0, 0, 0],
});
}
let time_series_table = db.get_time_series_table_mut(0).unwrap();
unsafe {
let result =
time_series_table.batch_write(more_data_points.as_ptr(), more_data_points.len());
assert!(result.is_ok());
assert_eq!(result.unwrap(), 10);
}
{
let time_series_table = db.get_time_series_table(0).unwrap();
let start_time = 1000000000000; let end_time = 1000000000000 + 1000000000; let interval_seconds = 60;
let avg_aggregation = "avg";
let result = time_series_table.query_pre_aggregated(
start_time,
end_time,
interval_seconds,
avg_aggregation,
);
assert!(result.is_ok());
let avg_records = result.unwrap();
assert!(!avg_records.is_empty());
println!("平均值预聚合查询结果: {:?}", avg_records);
}
println!("测试通过: 时序数据预聚合功能验证成功");
}