#![allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
use crate::metadata::logical_schema::{
LogicalDataType, LogicalField, LogicalSchema, LogicalTimestampUnit,
};
use crate::metadata::segments::{FileFormat, SegmentMeta};
use crate::metadata::table_metadata::{
TABLE_FORMAT_VERSION, TableKind, TableMeta, TimeBucket, TimeIndexSpec,
};
use crate::storage::{StorageError, TableLocation, layout};
use crate::transaction_log::{CommitError, LogAction, TransactionLogStore};
use chrono::{DateTime, TimeZone, Utc};
use tempfile::TempDir;
type TestResult = Result<(), Box<dyn std::error::Error>>;
fn create_test_log_store() -> (TempDir, TransactionLogStore) {
let tmp = TempDir::new().expect("create temp dir");
let location = TableLocation::local(tmp.path());
let store = TransactionLogStore::new(location);
(tmp, store)
}
fn sample_time_index_spec() -> TimeIndexSpec {
TimeIndexSpec {
timestamp_column: "ts".to_string(),
entity_columns: vec!["symbol".to_string()],
bucket: TimeBucket::Minutes(1),
timezone: None,
}
}
fn sample_table_meta() -> TableMeta {
let schema = LogicalSchema::new(vec![
LogicalField {
name: "ts".to_string(),
data_type: LogicalDataType::Timestamp {
unit: LogicalTimestampUnit::Micros,
timezone: None,
},
nullable: false,
},
LogicalField {
name: "symbol".to_string(),
data_type: LogicalDataType::Utf8,
nullable: false,
},
LogicalField {
name: "price".to_string(),
data_type: LogicalDataType::Float64,
nullable: true,
},
])
.expect("valid logical schema");
TableMeta::new_time_series_with_schema(sample_time_index_spec(), schema)
}
fn sample_segment(id: &str, ts_hour: u32) -> SegmentMeta {
SegmentMeta {
path: format!("data/{id}.parquet"),
format: FileFormat::Parquet,
ts_min: utc_datetime(2025, 1, 1, ts_hour, 0, 0),
ts_max: utc_datetime(2025, 1, 1, ts_hour + 1, 0, 0),
row_count: 1000,
file_size: None,
coverage_path: None,
}
}
fn utc_datetime(
year: i32,
month: u32,
day: u32,
hour: u32,
minute: u32,
second: u32,
) -> DateTime<Utc> {
Utc.with_ymd_and_hms(year, month, day, hour, minute, second)
.single()
.expect("valid UTC timestamp")
}
#[tokio::test]
async fn fresh_directory_returns_version_zero() -> TestResult {
let (_tmp, store) = create_test_log_store();
let version = store.load_current_version().await?;
assert_eq!(version, 0);
Ok(())
}
#[tokio::test]
async fn happy_path_commit_and_rebuild_table_state() -> TestResult {
let (tmp, store) = create_test_log_store();
let meta = sample_table_meta();
let seg1 = sample_segment("seg-001", 0);
let seg2 = sample_segment("seg-002", 1);
let v1 = store
.commit_with_expected_version(
0,
vec![
LogAction::UpdateTableMeta(meta.clone()),
LogAction::AddSegment(seg1.clone()),
],
)
.await?;
assert_eq!(v1, 1);
let commit_1_path = tmp.path().join(layout::commit_rel_path(1));
assert!(
commit_1_path.exists(),
"commit file for version 1 should exist"
);
let v2 = store
.commit_with_expected_version(v1, vec![LogAction::AddSegment(seg2.clone())])
.await?;
assert_eq!(v2, 2);
let commit_2_path = tmp.path().join(layout::commit_rel_path(2));
assert!(
commit_2_path.exists(),
"commit file for version 2 should exist"
);
let state = store.rebuild_table_state().await?;
assert_eq!(state.version, 2);
assert_eq!(state.segments.len(), 2);
assert!(state.segments.contains_key(&seg1.path));
assert!(state.segments.contains_key(&seg2.path));
match state.table_meta.kind() {
TableKind::TimeSeries(spec) => {
assert_eq!(spec.timestamp_column, "ts");
assert_eq!(spec.entity_columns, vec!["symbol".to_string()]);
}
TableKind::Generic => panic!("expected TimeSeries, got Generic"),
}
assert!(state.table_meta.logical_schema().is_some());
let schema = state
.table_meta
.logical_schema()
.expect("logical schema must be present");
assert_eq!(schema.columns().len(), 3);
Ok(())
}
#[tokio::test]
async fn sequential_commits_accumulate_segments() -> TestResult {
let (_tmp, store) = create_test_log_store();
let meta = sample_table_meta();
let v1 = store
.commit_with_expected_version(0, vec![LogAction::UpdateTableMeta(meta.clone())])
.await?;
assert_eq!(v1, 1);
let mut expected_version = v1;
for i in 1..=4 {
let seg = sample_segment(&format!("seg-{i:03}"), i as u32);
let v = store
.commit_with_expected_version(expected_version, vec![LogAction::AddSegment(seg)])
.await?;
assert_eq!(v, expected_version + 1);
expected_version = v;
}
let state = store.rebuild_table_state().await?;
assert_eq!(state.version, 5);
assert_eq!(state.segments.len(), 4);
Ok(())
}
#[tokio::test]
async fn remove_segment_removes_from_state() -> TestResult {
let (_tmp, store) = create_test_log_store();
let meta = sample_table_meta();
let seg1 = sample_segment("seg-to-keep", 0);
let seg2 = sample_segment("seg-to-remove", 1);
let v1 = store
.commit_with_expected_version(
0,
vec![
LogAction::UpdateTableMeta(meta),
LogAction::AddSegment(seg1.clone()),
LogAction::AddSegment(seg2.clone()),
],
)
.await?;
let v2 = store
.commit_with_expected_version(
v1,
vec![LogAction::RemoveSegment {
path: seg2.path.clone(),
}],
)
.await?;
let state = store.rebuild_table_state().await?;
assert_eq!(state.version, v2);
assert_eq!(state.segments.len(), 1);
assert!(state.segments.contains_key(&seg1.path));
assert!(!state.segments.contains_key(&seg2.path));
Ok(())
}
#[tokio::test]
async fn conflict_when_expected_version_is_stale() -> TestResult {
let (_tmp, store) = create_test_log_store();
let meta = sample_table_meta();
let expected_version = store.load_current_version().await?;
assert_eq!(expected_version, 0);
let v1 = store
.commit_with_expected_version(
expected_version,
vec![LogAction::UpdateTableMeta(meta.clone())],
)
.await?;
assert_eq!(v1, 1);
let result = store
.commit_with_expected_version(
expected_version,
vec![LogAction::AddSegment(sample_segment("seg-conflict", 0))],
)
.await;
match result {
Err(CommitError::Conflict {
expected, found, ..
}) => {
assert_eq!(expected, 0);
assert_eq!(found, 1);
}
other => panic!("expected Conflict error, got: {other:?}"),
}
let current = store.load_current_version().await?;
assert_eq!(current, 1);
Ok(())
}
#[tokio::test]
async fn conflict_on_subsequent_version() -> TestResult {
let (_tmp, store) = create_test_log_store();
let meta = sample_table_meta();
store
.commit_with_expected_version(0, vec![LogAction::UpdateTableMeta(meta)])
.await?;
store
.commit_with_expected_version(1, vec![LogAction::AddSegment(sample_segment("seg-1", 0))])
.await?;
let result = store
.commit_with_expected_version(1, vec![LogAction::AddSegment(sample_segment("seg-2", 1))])
.await;
match result {
Err(CommitError::Conflict {
expected, found, ..
}) => {
assert_eq!(expected, 1);
assert_eq!(found, 2);
}
other => panic!("expected Conflict error, got: {other:?}"),
}
Ok(())
}
#[tokio::test]
async fn corrupt_current_file_returns_corrupt_state() -> TestResult {
let (tmp, store) = create_test_log_store();
let log_dir = tmp.path().join(layout::log_rel_dir());
tokio::fs::create_dir_all(&log_dir).await?;
let current_path = tmp.path().join(layout::current_rel_path());
tokio::fs::write(¤t_path, "not-a-number").await?;
let result = store.load_current_version().await;
assert!(
matches!(result, Err(CommitError::CorruptState { .. })),
"expected CorruptState, got: {result:?}"
);
Ok(())
}
#[tokio::test]
async fn empty_current_file_returns_corrupt_state() -> TestResult {
let (tmp, store) = create_test_log_store();
let log_dir = tmp.path().join(layout::log_rel_dir());
tokio::fs::create_dir_all(&log_dir).await?;
let current_path = tmp.path().join(layout::current_rel_path());
tokio::fs::write(¤t_path, "").await?;
let result = store.load_current_version().await;
assert!(
matches!(result, Err(CommitError::CorruptState { .. })),
"expected CorruptState, got: {result:?}"
);
Ok(())
}
#[tokio::test]
async fn corrupt_commit_file_returns_corrupt_state() -> TestResult {
let (tmp, store) = create_test_log_store();
let meta = sample_table_meta();
store
.commit_with_expected_version(0, vec![LogAction::UpdateTableMeta(meta)])
.await?;
let commit_path = tmp.path().join(layout::commit_rel_path(1));
tokio::fs::write(&commit_path, "{ invalid json }}}").await?;
let result = store.load_commit(1).await;
assert!(
matches!(result, Err(CommitError::CorruptState { .. })),
"expected CorruptState, got: {result:?}"
);
let result = store.rebuild_table_state().await;
assert!(
matches!(result, Err(CommitError::CorruptState { .. })),
"expected CorruptState, got: {result:?}"
);
Ok(())
}
#[tokio::test]
async fn missing_commit_file_returns_storage_not_found() -> TestResult {
let (tmp, store) = create_test_log_store();
let meta = sample_table_meta();
store
.commit_with_expected_version(0, vec![LogAction::UpdateTableMeta(meta)])
.await?;
let commit_path = tmp.path().join(layout::commit_rel_path(1));
tokio::fs::remove_file(&commit_path).await?;
let result = store.load_commit(1).await;
match result {
Err(CommitError::Storage {
source: StorageError::NotFound { .. },
}) => {}
other => panic!("expected Storage(NotFound), got: {other:?}"),
}
Ok(())
}
#[tokio::test]
async fn leftover_tmp_files_are_ignored() -> TestResult {
let (tmp, store) = create_test_log_store();
let meta = sample_table_meta();
let seg = sample_segment("seg-1", 0);
store
.commit_with_expected_version(
0,
vec![
LogAction::UpdateTableMeta(meta.clone()),
LogAction::AddSegment(seg.clone()),
],
)
.await?;
let log_dir = tmp.path().join(layout::log_rel_dir());
tokio::fs::write(log_dir.join("0000000002.json.tmp"), b"garbage").await?;
tokio::fs::write(log_dir.join(".tmp_random_file"), b"more garbage").await?;
tokio::fs::write(log_dir.join("temp_commit.tmp"), b"even more garbage").await?;
let state = store.rebuild_table_state().await?;
assert_eq!(state.version, 1);
assert_eq!(state.segments.len(), 1);
assert!(log_dir.join("0000000002.json.tmp").exists());
Ok(())
}
#[tokio::test]
async fn missing_intermediate_commit_fails_rebuild() -> TestResult {
let (tmp, store) = create_test_log_store();
let meta = sample_table_meta();
store
.commit_with_expected_version(0, vec![LogAction::UpdateTableMeta(meta.clone())])
.await?;
store
.commit_with_expected_version(1, vec![LogAction::AddSegment(sample_segment("seg-1", 0))])
.await?;
let commit_1_path = tmp.path().join(layout::commit_rel_path(1));
tokio::fs::remove_file(&commit_1_path).await?;
let result = store.rebuild_table_state().await;
match result {
Err(CommitError::Storage {
source: StorageError::NotFound { .. },
}) => {}
other => panic!("expected Storage(NotFound), got: {other:?}"),
}
Ok(())
}
#[tokio::test]
async fn rebuild_on_empty_table_returns_corrupt_state() -> TestResult {
let (_tmp, store) = create_test_log_store();
let result = store.rebuild_table_state().await;
assert!(
matches!(result, Err(CommitError::CorruptState { .. })),
"expected CorruptState for empty table, got: {result:?}"
);
Ok(())
}
#[tokio::test]
async fn rebuild_without_table_meta_returns_corrupt_state() -> TestResult {
let (_tmp, store) = create_test_log_store();
store
.commit_with_expected_version(0, vec![LogAction::AddSegment(sample_segment("seg-1", 0))])
.await?;
let result = store.rebuild_table_state().await;
assert!(
matches!(result, Err(CommitError::CorruptState { .. })),
"expected CorruptState for missing TableMeta, got: {result:?}"
);
Ok(())
}
#[tokio::test]
async fn update_table_meta_last_one_wins() -> TestResult {
let (_tmp, store) = create_test_log_store();
let meta1 = TableMeta::new_time_series(TimeIndexSpec {
timestamp_column: "ts".to_string(),
entity_columns: vec![],
bucket: TimeBucket::Minutes(1),
timezone: None,
});
let meta2 = TableMeta::new_time_series(TimeIndexSpec {
timestamp_column: "event_time".to_string(), entity_columns: vec!["user_id".to_string()],
bucket: TimeBucket::Hours(1),
timezone: Some("UTC".to_string()),
});
store
.commit_with_expected_version(0, vec![LogAction::UpdateTableMeta(meta1)])
.await?;
store
.commit_with_expected_version(1, vec![LogAction::UpdateTableMeta(meta2.clone())])
.await?;
let state = store.rebuild_table_state().await?;
match state.table_meta.kind() {
TableKind::TimeSeries(spec) => {
assert_eq!(spec.timestamp_column, "event_time");
assert_eq!(spec.entity_columns, vec!["user_id".to_string()]);
assert_eq!(spec.bucket, TimeBucket::Hours(1));
}
_ => panic!("expected TimeSeries"),
}
assert_eq!(state.table_meta.format_version(), TABLE_FORMAT_VERSION);
Ok(())
}
#[tokio::test]
async fn duplicate_live_segment_path_is_corrupt_state() -> TestResult {
let (_tmp, store) = create_test_log_store();
let meta = sample_table_meta();
let seg = sample_segment("seg-001", 0);
let mut duplicate = sample_segment("seg-002", 1);
duplicate.path = seg.path.clone();
store
.commit_with_expected_version(
0,
vec![
LogAction::UpdateTableMeta(meta),
LogAction::AddSegment(seg.clone()),
],
)
.await?;
store
.commit_with_expected_version(1, vec![LogAction::AddSegment(duplicate)])
.await?;
let err = store
.rebuild_table_state()
.await
.expect_err("duplicate live path must be corrupt");
assert!(matches!(
err,
CommitError::CorruptState { ref msg, .. } if msg.contains(&seg.path)
));
Ok(())
}
#[tokio::test]
async fn remove_nonexistent_segment_is_noop() -> TestResult {
let (_tmp, store) = create_test_log_store();
let meta = sample_table_meta();
let seg = sample_segment("seg-exists", 0);
store
.commit_with_expected_version(
0,
vec![
LogAction::UpdateTableMeta(meta),
LogAction::AddSegment(seg.clone()),
],
)
.await?;
store
.commit_with_expected_version(
1,
vec![LogAction::RemoveSegment {
path: "data/does-not-exist.parquet".to_string(),
}],
)
.await?;
let state = store.rebuild_table_state().await?;
assert_eq!(state.segments.len(), 1);
assert!(state.segments.contains_key(&seg.path));
Ok(())
}
#[tokio::test]
async fn table_coverage_pointer_is_replayed() -> TestResult {
let (_tmp, store) = create_test_log_store();
let meta = sample_table_meta();
let coverage_bucket = TimeBucket::Minutes(5);
let coverage_path = "coverage/0000000002.bitmap".to_string();
let v1 = store
.commit_with_expected_version(0, vec![LogAction::UpdateTableMeta(meta)])
.await?;
let v2 = store
.commit_with_expected_version(
v1,
vec![LogAction::UpdateTableCoverage {
bucket_spec: coverage_bucket.clone(),
coverage_path: coverage_path.clone(),
}],
)
.await?;
let state = store.rebuild_table_state().await?;
assert_eq!(state.version, v2);
let pointer = state
.table_coverage
.as_ref()
.expect("table coverage pointer should be present");
assert_eq!(pointer.bucket_spec, coverage_bucket);
assert_eq!(pointer.coverage_path, coverage_path);
assert_eq!(pointer.version, v2);
Ok(())
}
#[tokio::test]
async fn table_coverage_last_one_wins() -> TestResult {
let (_tmp, store) = create_test_log_store();
let meta = sample_table_meta();
let coverage_bucket_v1 = TimeBucket::Minutes(5);
let coverage_bucket_v2 = TimeBucket::Hours(1);
let coverage_path_v1 = "coverage/0000000002.bitmap".to_string();
let coverage_path_v2 = "coverage/0000000003.bitmap".to_string();
let v1 = store
.commit_with_expected_version(0, vec![LogAction::UpdateTableMeta(meta)])
.await?;
let v2 = store
.commit_with_expected_version(
v1,
vec![LogAction::UpdateTableCoverage {
bucket_spec: coverage_bucket_v1.clone(),
coverage_path: coverage_path_v1.clone(),
}],
)
.await?;
let v3 = store
.commit_with_expected_version(
v2,
vec![LogAction::UpdateTableCoverage {
bucket_spec: coverage_bucket_v2.clone(),
coverage_path: coverage_path_v2.clone(),
}],
)
.await?;
let state = store.rebuild_table_state().await?;
assert_eq!(state.version, v3);
let pointer = state
.table_coverage
.as_ref()
.expect("table coverage pointer should be present");
assert_eq!(pointer.bucket_spec, coverage_bucket_v2);
assert_eq!(pointer.coverage_path, coverage_path_v2);
assert_eq!(pointer.version, v3);
Ok(())
}
#[tokio::test]
async fn table_coverage_is_none_when_not_committed() -> TestResult {
let (_tmp, store) = create_test_log_store();
let meta = sample_table_meta();
let seg = sample_segment("seg-without-coverage", 0);
let v1 = store
.commit_with_expected_version(0, vec![LogAction::UpdateTableMeta(meta)])
.await?;
store
.commit_with_expected_version(v1, vec![LogAction::AddSegment(seg)])
.await?;
let state = store.rebuild_table_state().await?;
assert_eq!(state.version, 2);
assert!(state.table_coverage.is_none());
Ok(())
}
#[tokio::test]
async fn table_coverage_rebuilds_with_segment_coverage_paths() -> TestResult {
let (_tmp, store) = create_test_log_store();
let meta = sample_table_meta();
let coverage_bucket = TimeBucket::Minutes(5);
let segment_cov_path = "coverage/seg-001.roar".to_string();
let snapshot_cov_path = "coverage/table/0000000002.roar".to_string();
let mut seg = sample_segment("seg-001", 0);
seg.coverage_path = Some(segment_cov_path.clone());
let v1 = store
.commit_with_expected_version(0, vec![LogAction::UpdateTableMeta(meta)])
.await?;
let v2 = store
.commit_with_expected_version(v1, vec![LogAction::AddSegment(seg.clone())])
.await?;
let v3 = store
.commit_with_expected_version(
v2,
vec![LogAction::UpdateTableCoverage {
bucket_spec: coverage_bucket.clone(),
coverage_path: snapshot_cov_path.clone(),
}],
)
.await?;
let state = store.rebuild_table_state().await?;
assert_eq!(state.version, v3);
let rebuilt_seg = state
.segments
.get(&seg.path)
.expect("segment present after rebuild");
assert_eq!(
rebuilt_seg.coverage_path.as_deref(),
Some(segment_cov_path.as_str())
);
let pointer = state
.table_coverage
.as_ref()
.expect("table coverage pointer should be present");
assert_eq!(pointer.bucket_spec, coverage_bucket);
assert_eq!(pointer.coverage_path, snapshot_cov_path);
assert_eq!(pointer.version, v3);
Ok(())
}