use std::sync::Arc;
use lance_core::Error;
use lance_core::deepsize::DeepSizeOf;
use lance_core::error::Result;
use roaring::RoaringBitmap;
use serde::{Deserialize, Serialize};
use object_store::path::Path;
use super::DataFile;
use crate::format::pb;
pub const TOMBSTONE_FIELD_ID: i32 = -2;
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(into = "OverlayCoverageBytes", try_from = "OverlayCoverageBytes")]
pub enum OverlayCoverage {
Shared(Arc<RoaringBitmap>),
PerField(Vec<Arc<RoaringBitmap>>),
}
#[derive(Debug, Clone, Serialize, Deserialize)]
enum OverlayCoverageBytes {
Shared(Vec<u8>),
PerField(Vec<Vec<u8>>),
}
fn deserialize_roaring(bytes: &[u8], path: &Path) -> Result<RoaringBitmap> {
RoaringBitmap::deserialize_from(bytes).map_err(|e| {
Error::corrupt_file(
path.clone(),
format!("failed to deserialize overlay coverage bitmap: {e}"),
)
})
}
fn serialize_roaring(bitmap: &RoaringBitmap) -> Vec<u8> {
let mut bytes = Vec::with_capacity(bitmap.serialized_size());
bitmap.serialize_into(&mut bytes).unwrap();
bytes
}
impl From<OverlayCoverage> for OverlayCoverageBytes {
fn from(coverage: OverlayCoverage) -> Self {
match coverage {
OverlayCoverage::Shared(bitmap) => Self::Shared(serialize_roaring(&bitmap)),
OverlayCoverage::PerField(bitmaps) => {
Self::PerField(bitmaps.iter().map(|b| serialize_roaring(b)).collect())
}
}
}
}
impl TryFrom<OverlayCoverageBytes> for OverlayCoverage {
type Error = Error;
fn try_from(bytes: OverlayCoverageBytes) -> Result<Self> {
let path = Path::default();
Ok(match bytes {
OverlayCoverageBytes::Shared(b) => {
Self::Shared(Arc::new(deserialize_roaring(&b, &path)?))
}
OverlayCoverageBytes::PerField(bs) => Self::PerField(
bs.iter()
.map(|b| deserialize_roaring(b, &path).map(Arc::new))
.collect::<Result<_>>()?,
),
})
}
}
impl DeepSizeOf for OverlayCoverage {
fn deep_size_of_children(&self, context: &mut lance_core::deepsize::Context) -> usize {
let bitmap_heap = |bitmap: &Arc<RoaringBitmap>,
context: &mut lance_core::deepsize::Context| {
if context.mark_seen(Arc::as_ptr(bitmap) as usize) {
std::mem::size_of::<RoaringBitmap>() + bitmap.serialized_size()
} else {
0
}
};
match self {
Self::Shared(bitmap) => bitmap_heap(bitmap, context),
Self::PerField(bitmaps) => {
bitmaps.capacity() * std::mem::size_of::<Arc<RoaringBitmap>>()
+ bitmaps
.iter()
.map(|b| bitmap_heap(b, context))
.sum::<usize>()
}
}
}
}
impl OverlayCoverage {
pub fn dense(bitmap: RoaringBitmap) -> Self {
Self::Shared(Arc::new(bitmap))
}
pub fn sparse(bitmaps: Vec<RoaringBitmap>) -> Self {
Self::PerField(bitmaps.into_iter().map(Arc::new).collect())
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, DeepSizeOf)]
pub struct DataOverlayFile {
pub data_file: DataFile,
pub coverage: OverlayCoverage,
pub committed_version: u64,
}
impl DataOverlayFile {
pub fn coverage_for_field(&self, field_pos: usize) -> Result<Arc<RoaringBitmap>> {
match &self.coverage {
OverlayCoverage::Shared(bitmap) => Ok(bitmap.clone()),
OverlayCoverage::PerField(bitmaps) => {
bitmaps.get(field_pos).cloned().ok_or_else(|| {
Error::invalid_input(format!(
"overlay per-field coverage has {} bitmaps but field position {} was requested",
bitmaps.len(),
field_pos
))
})
}
}
}
}
pub fn sort_overlays_newest_last(overlays: &mut [DataOverlayFile]) {
overlays.sort_by_key(|overlay| overlay.committed_version);
}
pub fn verify_overlays_newest_last(overlays: &[DataOverlayFile]) -> Result<()> {
for pair in overlays.windows(2) {
if pair[0].committed_version > pair[1].committed_version {
return Err(Error::invalid_input(format!(
"overlay files must be stored newest-last, but committed_version {} precedes {}",
pair[0].committed_version, pair[1].committed_version
)));
}
}
Ok(())
}
pub fn tombstone_overlay_fields(overlays: &mut Vec<DataOverlayFile>, fields: &[u32]) {
for overlay in overlays.iter_mut() {
let tombstoned: Vec<i32> = overlay
.data_file
.fields
.iter()
.map(|&field| {
if field >= 0 && fields.contains(&(field as u32)) {
TOMBSTONE_FIELD_ID
} else {
field
}
})
.collect();
overlay.data_file.fields = tombstoned.into();
}
overlays.retain(|overlay| {
overlay
.data_file
.fields
.iter()
.any(|&field| field != TOMBSTONE_FIELD_ID)
});
}
impl From<&DataOverlayFile> for pb::DataOverlayFile {
fn from(overlay: &DataOverlayFile) -> Self {
let coverage = match &overlay.coverage {
OverlayCoverage::Shared(bitmap) => {
pb::data_overlay_file::Coverage::SharedOffsetBitmap(serialize_roaring(bitmap))
}
OverlayCoverage::PerField(bitmaps) => {
pb::data_overlay_file::Coverage::FieldCoverage(pb::FieldCoverage {
offset_bitmaps: bitmaps.iter().map(|b| serialize_roaring(b)).collect(),
})
}
};
Self {
data_file: Some(pb::DataFile::from(&overlay.data_file)),
coverage: Some(coverage),
committed_version: overlay.committed_version,
}
}
}
impl TryFrom<pb::DataOverlayFile> for DataOverlayFile {
type Error = Error;
fn try_from(proto: pb::DataOverlayFile) -> Result<Self> {
let data_file = proto
.data_file
.ok_or_else(|| Error::invalid_input("DataOverlayFile is missing its data_file"))?;
let path = Path::from(data_file.path.as_str());
let coverage = match proto.coverage {
Some(pb::data_overlay_file::Coverage::SharedOffsetBitmap(bytes)) => {
OverlayCoverage::Shared(Arc::new(deserialize_roaring(&bytes, &path)?))
}
Some(pb::data_overlay_file::Coverage::FieldCoverage(fc)) => OverlayCoverage::PerField(
fc.offset_bitmaps
.iter()
.map(|b| deserialize_roaring(b, &path).map(Arc::new))
.collect::<Result<_>>()?,
),
None => {
return Err(Error::invalid_input(
"DataOverlayFile is missing its coverage",
));
}
};
Ok(Self {
data_file: DataFile::try_from(data_file)?,
coverage,
committed_version: proto.committed_version,
})
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_data_overlay_missing_fields_error() {
let no_coverage = pb::DataOverlayFile {
data_file: Some(pb::DataFile::from(&DataFile::new_legacy_from_fields(
"overlay.lance",
vec![3],
None,
))),
coverage: None,
committed_version: 1,
};
let err = DataOverlayFile::try_from(no_coverage).unwrap_err();
assert!(err.to_string().contains("missing its coverage"), "{err}");
let no_data_file = pb::DataOverlayFile {
data_file: None,
coverage: Some(pb::data_overlay_file::Coverage::SharedOffsetBitmap(
serialize_roaring(&RoaringBitmap::from_iter([0u32])),
)),
committed_version: 1,
};
let err = DataOverlayFile::try_from(no_data_file).unwrap_err();
assert!(err.to_string().contains("missing its data_file"), "{err}");
}
#[test]
fn test_overlay_coverage_serde_json_roundtrip() {
for coverage in [
OverlayCoverage::dense(RoaringBitmap::from_iter([1u32, 5, 100])),
OverlayCoverage::dense(RoaringBitmap::new()),
OverlayCoverage::sparse(vec![
RoaringBitmap::from_iter([2u32, 3]),
RoaringBitmap::new(),
]),
OverlayCoverage::sparse(vec![]),
] {
let json = serde_json::to_string(&coverage).unwrap();
let back: OverlayCoverage = serde_json::from_str(&json).unwrap();
assert_eq!(back, coverage);
}
}
#[test]
fn test_tombstone_overlay_fields() {
let mut overlays = vec![
DataOverlayFile {
data_file: DataFile::new_legacy_from_fields("a.lance", vec![3, 5], None),
coverage: OverlayCoverage::sparse(vec![
RoaringBitmap::from_iter([0u32]),
RoaringBitmap::from_iter([1u32]),
]),
committed_version: 1,
},
DataOverlayFile {
data_file: DataFile::new_legacy_from_fields("b.lance", vec![5], None),
coverage: OverlayCoverage::dense(RoaringBitmap::from_iter([0u32])),
committed_version: 1,
},
DataOverlayFile {
data_file: DataFile::new_legacy_from_fields("c.lance", vec![7], None),
coverage: OverlayCoverage::dense(RoaringBitmap::from_iter([0u32])),
committed_version: 1,
},
];
tombstone_overlay_fields(&mut overlays, &[5]);
assert_eq!(overlays.len(), 2);
assert_eq!(
overlays[0].data_file.fields.as_ref(),
&[3, TOMBSTONE_FIELD_ID]
);
assert_eq!(overlays[1].data_file.fields.as_ref(), &[7]);
}
#[test]
fn test_verify_overlays_newest_last() {
let mk = |version: u64| DataOverlayFile {
data_file: DataFile::new_legacy_from_fields("o.lance", vec![3], None),
coverage: OverlayCoverage::dense(RoaringBitmap::from_iter([0u32])),
committed_version: version,
};
assert!(verify_overlays_newest_last(&[]).is_ok());
assert!(verify_overlays_newest_last(&[mk(1), mk(2), mk(2), mk(5)]).is_ok());
let err = verify_overlays_newest_last(&[mk(2), mk(1)]).unwrap_err();
assert!(err.to_string().contains("newest-last"), "{err}");
}
#[test]
fn test_coverage_for_field_out_of_bounds() {
let overlay = DataOverlayFile {
data_file: DataFile::new_legacy_from_fields("o.lance", vec![2, 4], None),
coverage: OverlayCoverage::sparse(vec![
RoaringBitmap::from_iter([1u32]),
RoaringBitmap::from_iter([2u32]),
]),
committed_version: 1,
};
assert!(overlay.coverage_for_field(0).is_ok());
assert!(overlay.coverage_for_field(1).is_ok());
let err = overlay.coverage_for_field(5).unwrap_err();
assert!(err.to_string().contains("field position"), "{err}");
}
}