mod api;
mod hnsw;
pub mod index;
mod manifest;
pub mod memtable;
pub mod observer;
pub mod scanner;
pub mod sharding;
#[cfg(test)]
pub(crate) mod test_util;
pub mod util;
mod wal;
pub mod write;
use std::sync::Arc;
use arrow_schema::{DataType, Field as ArrowField, Schema as ArrowSchema};
pub const TOMBSTONE: &str = "_tombstone";
pub fn tombstone_field() -> ArrowField {
ArrowField::new(TOMBSTONE, DataType::Boolean, false)
}
pub fn relax_non_pk_nullability(
logical_schema: &ArrowSchema,
pk_columns: &[String],
) -> Arc<ArrowSchema> {
let fields: Vec<ArrowField> = logical_schema
.fields()
.iter()
.map(|field| {
let keep = field.is_nullable()
|| field.name() == TOMBSTONE
|| pk_columns.iter().any(|c| c == field.name());
let field = field.as_ref().clone();
if keep {
field
} else {
field.with_nullable(true)
}
})
.collect();
Arc::new(ArrowSchema::new_with_metadata(
fields,
logical_schema.metadata().clone(),
))
}
pub fn schema_with_tombstone(base: &ArrowSchema) -> Arc<ArrowSchema> {
if base.column_with_name(TOMBSTONE).is_some() {
return Arc::new(base.clone());
}
let mut fields: Vec<ArrowField> = base.fields().iter().map(|f| f.as_ref().clone()).collect();
fields.push(tombstone_field());
Arc::new(ArrowSchema::new_with_metadata(
fields,
base.metadata().clone(),
))
}
pub use api::{DatasetMemWalExt, InitializeMemWalBuilder, validate_maintained_indexes};
pub use index::MemIndexKind;
pub use manifest::ShardManifestStore;
pub use memtable::scanner::MemTableScanner;
pub use scanner::{LsmDataSource, LsmGeneration, LsmScanner, ShardSnapshot};
pub use sharding::{
evaluate_sharding_spec, evaluate_sharding_spec_with_embedded_columns,
evaluate_sharding_spec_with_source_columns,
};
pub use wal::{BatchDurableWatcher, WalAppendResult, WalAppender, WalReadEntry, WalTailer};
pub use write::SealFence;
pub use write::ShardWriter;
pub use write::ShardWriterConfig;
pub use write::WriteResult;
#[cfg(test)]
mod tests {
use super::*;
use arrow_schema::Fields;
fn logical() -> ArrowSchema {
ArrowSchema::new(vec![
ArrowField::new("id", DataType::Int32, false),
ArrowField::new("count", DataType::Int64, false),
ArrowField::new("note", DataType::Utf8, true),
])
}
#[test]
fn relax_widens_every_non_pk_field_and_leaves_the_key_alone() {
let relaxed = relax_non_pk_nullability(&logical(), &["id".to_string()]);
assert!(
!relaxed.field(0).is_nullable(),
"the primary key stays strict"
);
assert!(
relaxed.field(1).is_nullable(),
"`count` must accept a tombstone null"
);
assert!(
relaxed.field(2).is_nullable(),
"already-nullable is untouched"
);
}
#[test]
fn relax_leaves_nested_fields_exactly_as_declared() {
let item = Arc::new(ArrowField::new("item", DataType::Float32, false));
let child = ArrowField::new("a", DataType::Int32, false);
let schema = ArrowSchema::new(vec![
ArrowField::new("id", DataType::Int32, false),
ArrowField::new("vector", DataType::FixedSizeList(item, 4), false),
ArrowField::new("s", DataType::Struct(Fields::from(vec![child])), false),
]);
let relaxed = relax_non_pk_nullability(&schema, &["id".to_string()]);
assert!(relaxed.field(1).is_nullable());
match relaxed.field(1).data_type() {
DataType::FixedSizeList(f, _) => assert!(!f.is_nullable(), "item field untouched"),
other => panic!("expected FixedSizeList, got {other:?}"),
}
match relaxed.field(2).data_type() {
DataType::Struct(fields) => assert!(!fields[0].is_nullable(), "child field untouched"),
other => panic!("expected Struct, got {other:?}"),
}
}
#[test]
fn relax_keeps_tombstone_non_nullable_and_is_idempotent() {
let pk = ["id".to_string()];
let once = relax_non_pk_nullability(&schema_with_tombstone(&logical()), &pk);
let twice = relax_non_pk_nullability(&once, &pk);
let tombstone = once.field_with_name(TOMBSTONE).unwrap();
assert!(
!tombstone.is_nullable(),
"the write path always populates _tombstone"
);
assert_eq!(once, twice);
}
#[test]
fn relax_preserves_schema_and_field_metadata() {
let marked = ArrowField::new("count", DataType::Int64, false)
.with_metadata([("k".to_string(), "v".to_string())].into());
let schema = ArrowSchema::new_with_metadata(
vec![ArrowField::new("id", DataType::Int32, false), marked],
[("s".to_string(), "m".to_string())].into(),
);
let relaxed = relax_non_pk_nullability(&schema, &["id".to_string()]);
assert_eq!(relaxed.metadata().get("s").map(String::as_str), Some("m"));
assert_eq!(
relaxed.field(1).metadata().get("k").map(String::as_str),
Some("v")
);
}
}