pub mod actions;
pub mod log_store;
pub mod segments;
pub mod table_state;
#[cfg(test)]
mod log_integration_tests;
pub use crate::metadata::table_metadata::{
TableKind, TableMeta, TableMetaDelta, TimeBucket, TimeIndexSpec,
};
pub use actions::{Commit, LogAction};
pub use log_store::TransactionLogStore;
pub use segments::{FileFormat, SegmentMeta};
pub use table_state::TableState;
use snafu::{Backtrace, prelude::*};
use crate::storage::StorageError;
#[derive(Debug, Snafu)]
pub enum CommitError {
#[snafu(display("Commit conflict: expected version {expected}, but CURRENT is {found}"))]
Conflict {
expected: u64,
found: u64,
backtrace: Backtrace,
},
#[snafu(display("Storage error while accessing commit log: {source}"))]
Storage {
#[snafu(backtrace)]
source: StorageError,
},
#[snafu(display("Corrupt log state: {msg}"))]
CorruptState {
msg: String,
backtrace: Backtrace,
},
}
#[cfg(test)]
mod tests {
use crate::metadata::logical_schema::{
LogicalDataType, LogicalField, LogicalSchema, LogicalSchemaError, LogicalTimestampUnit,
};
use crate::metadata::table_metadata::TABLE_FORMAT_VERSION;
use crate::transaction_log::*;
use chrono::{DateTime, TimeZone, Utc};
use serde_json;
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")
}
#[test]
fn commit_json_roundtrip() {
let ts0 = utc_datetime(2025, 1, 1, 0, 0, 0);
let ts1 = utc_datetime(2025, 1, 1, 1, 0, 0);
let time_index = TimeIndexSpec {
timestamp_column: "ts".to_string(),
entity_columns: vec!["symbol".to_string()],
bucket: TimeBucket::Minutes(60),
timezone: Some("UTC".to_string()),
};
let table_meta = TableMeta {
kind: TableKind::TimeSeries(time_index),
logical_schema: Some(
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,
},
])
.expect("valid logical schema"),
),
created_at: ts0,
format_version: TABLE_FORMAT_VERSION,
entity_identity: None,
};
let seg_meta = SegmentMeta {
path: "data/nvda_1h_0001.parquet".to_string(),
format: FileFormat::Parquet,
ts_min: ts0,
ts_max: ts1,
row_count: 1024,
file_size: None,
coverage_path: None,
};
let commit = Commit {
version: 1,
base_version: 0,
timestamp: ts1,
actions: vec![
LogAction::UpdateTableMeta(table_meta),
LogAction::AddSegment(seg_meta),
],
};
let json = serde_json::to_string_pretty(&commit).expect("serialize commit");
let decoded: Commit = serde_json::from_str(&json).expect("deserialize commit");
assert_eq!(commit, decoded);
}
#[test]
fn logical_schema_rejects_duplicate_columns() {
let dup = LogicalSchema::new(vec![
LogicalField {
name: "ts".to_string(),
data_type: LogicalDataType::Timestamp {
unit: LogicalTimestampUnit::Micros,
timezone: None,
},
nullable: false,
},
LogicalField {
name: "ts".to_string(),
data_type: LogicalDataType::Timestamp {
unit: LogicalTimestampUnit::Micros,
timezone: None,
},
nullable: false,
},
]);
let err = dup.expect_err("duplicate columns should be rejected");
assert!(matches!(err, LogicalSchemaError::DuplicateColumn { column } if column == "ts"));
}
#[test]
fn time_index_spec_defaults() {
let json = r#"{
"timestamp_column": "ts",
"bucket": { "Hours": 1 }
}"#;
let spec: TimeIndexSpec = serde_json::from_str(json).expect("deserialize");
assert_eq!(spec.timestamp_column, "ts");
assert_eq!(spec.entity_columns, Vec::<String>::new()); assert_eq!(spec.bucket, TimeBucket::Hours(1));
assert_eq!(spec.timezone, None); }
#[test]
fn time_index_spec_skips_none_timezone_on_serialize() {
let spec = TimeIndexSpec {
timestamp_column: "ts".to_string(),
entity_columns: vec![],
bucket: TimeBucket::Seconds(30),
timezone: None,
};
let json = serde_json::to_string(&spec).expect("serialize");
assert!(!json.contains("timezone"));
}
#[test]
fn logical_column_nullable_requires_explicit_value() {
let json = r#"{ "name": "price", "data_type": "Float64" }"#;
let err = serde_json::from_str::<LogicalField>(json).unwrap_err();
assert!(
err.to_string().contains("missing field `nullable`"),
"unexpected error: {err}"
);
}
#[test]
fn table_kind_generic_roundtrip() {
let kind = TableKind::Generic;
let json = serde_json::to_string(&kind).expect("serialize");
let decoded: TableKind = serde_json::from_str(&json).expect("deserialize");
assert_eq!(kind, decoded);
assert_eq!(json, r#""Generic""#);
}
#[test]
fn all_time_bucket_variants_roundtrip() {
let buckets = vec![
TimeBucket::Seconds(15),
TimeBucket::Minutes(5),
TimeBucket::Hours(24),
TimeBucket::Days(7),
];
for bucket in buckets {
let json = serde_json::to_string(&bucket).expect("serialize");
let decoded: TimeBucket = serde_json::from_str(&json).expect("deserialize");
assert_eq!(bucket, decoded);
}
}
#[test]
fn file_format_serializes_lowercase() {
let format = FileFormat::Parquet;
let json = serde_json::to_string(&format).expect("serialize");
assert_eq!(json, r#""parquet""#);
let decoded: FileFormat = serde_json::from_str(&json).expect("deserialize");
assert_eq!(format, decoded);
}
#[test]
fn file_format_default_is_parquet() {
assert_eq!(FileFormat::default(), FileFormat::Parquet);
}
#[test]
fn remove_segment_action_roundtrip() {
let action = LogAction::RemoveSegment {
path: "data/seg-to-remove.parquet".to_string(),
};
let json = serde_json::to_string(&action).expect("serialize");
let decoded: LogAction = serde_json::from_str(&json).expect("deserialize");
assert_eq!(action, decoded);
}
#[test]
fn commit_with_empty_actions() {
let ts = utc_datetime(2025, 6, 15, 12, 0, 0);
let commit = Commit {
version: 1,
base_version: 0,
timestamp: ts,
actions: vec![],
};
let json = serde_json::to_string(&commit).expect("serialize");
let decoded: Commit = serde_json::from_str(&json).expect("deserialize");
assert_eq!(commit, decoded);
assert!(decoded.actions.is_empty());
}
}