1pub mod staleness;
37
38use std::sync::Arc;
39
40use lance_core::Error;
41use lance_core::deepsize::DeepSizeOf;
42use lance_core::error::Result;
43use roaring::RoaringBitmap;
44use serde::{Deserialize, Serialize};
45
46use object_store::path::Path;
47
48use super::DataFile;
49use crate::format::pb;
50
51pub const TOMBSTONE_FIELD_ID: i32 = -2;
55
56#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
64#[serde(into = "OverlayCoverageBytes", try_from = "OverlayCoverageBytes")]
65pub enum OverlayCoverage {
66 Shared(Arc<RoaringBitmap>),
70 PerField(Vec<Arc<RoaringBitmap>>),
74}
75
76#[derive(Debug, Clone, Serialize, Deserialize)]
79enum OverlayCoverageBytes {
80 Shared(Vec<u8>),
81 PerField(Vec<Vec<u8>>),
82}
83
84fn deserialize_roaring(bytes: &[u8], path: &Path) -> Result<RoaringBitmap> {
89 RoaringBitmap::deserialize_from(bytes).map_err(|e| {
90 Error::corrupt_file(
91 path.clone(),
92 format!("failed to deserialize overlay coverage bitmap: {e}"),
93 )
94 })
95}
96
97fn serialize_roaring(bitmap: &RoaringBitmap) -> Vec<u8> {
98 let mut bitmap = bitmap.clone();
99 bitmap.optimize();
100 let mut bytes = Vec::with_capacity(bitmap.serialized_size());
101 bitmap.serialize_into(&mut bytes).unwrap();
103 bytes
104}
105
106impl From<OverlayCoverage> for OverlayCoverageBytes {
107 fn from(coverage: OverlayCoverage) -> Self {
108 match coverage {
109 OverlayCoverage::Shared(bitmap) => Self::Shared(serialize_roaring(&bitmap)),
110 OverlayCoverage::PerField(bitmaps) => {
111 Self::PerField(bitmaps.iter().map(|b| serialize_roaring(b)).collect())
112 }
113 }
114 }
115}
116
117impl TryFrom<OverlayCoverageBytes> for OverlayCoverage {
118 type Error = Error;
119
120 fn try_from(bytes: OverlayCoverageBytes) -> Result<Self> {
121 let path = Path::default();
124 Ok(match bytes {
125 OverlayCoverageBytes::Shared(b) => {
126 Self::Shared(Arc::new(deserialize_roaring(&b, &path)?))
127 }
128 OverlayCoverageBytes::PerField(bs) => Self::PerField(
129 bs.iter()
130 .map(|b| deserialize_roaring(b, &path).map(Arc::new))
131 .collect::<Result<_>>()?,
132 ),
133 })
134 }
135}
136
137impl DeepSizeOf for OverlayCoverage {
138 fn deep_size_of_children(&self, context: &mut lance_core::deepsize::Context) -> usize {
139 let bitmap_heap = |bitmap: &Arc<RoaringBitmap>,
145 context: &mut lance_core::deepsize::Context| {
146 if context.mark_seen(Arc::as_ptr(bitmap) as usize) {
147 std::mem::size_of::<RoaringBitmap>() + bitmap.serialized_size()
148 } else {
149 0
150 }
151 };
152 match self {
153 Self::Shared(bitmap) => bitmap_heap(bitmap, context),
154 Self::PerField(bitmaps) => {
155 bitmaps.capacity() * std::mem::size_of::<Arc<RoaringBitmap>>()
156 + bitmaps
157 .iter()
158 .map(|b| bitmap_heap(b, context))
159 .sum::<usize>()
160 }
161 }
162 }
163}
164
165impl OverlayCoverage {
166 pub fn dense(bitmap: RoaringBitmap) -> Self {
168 Self::Shared(Arc::new(bitmap))
169 }
170
171 pub fn sparse(bitmaps: Vec<RoaringBitmap>) -> Self {
173 Self::PerField(bitmaps.into_iter().map(Arc::new).collect())
174 }
175}
176
177#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, DeepSizeOf)]
181pub struct DataOverlayFile {
182 pub data_file: DataFile,
184 pub coverage: OverlayCoverage,
186 pub committed_version: u64,
190}
191
192impl DataOverlayFile {
193 pub fn coverage_for_field(&self, field_pos: usize) -> Result<Arc<RoaringBitmap>> {
200 match &self.coverage {
201 OverlayCoverage::Shared(bitmap) => Ok(bitmap.clone()),
202 OverlayCoverage::PerField(bitmaps) => {
203 bitmaps.get(field_pos).cloned().ok_or_else(|| {
204 Error::invalid_input(format!(
205 "overlay per-field coverage has {} bitmaps but field position {} was requested",
206 bitmaps.len(),
207 field_pos
208 ))
209 })
210 }
211 }
212 }
213}
214
215pub fn sort_overlays_newest_last(overlays: &mut [DataOverlayFile]) {
220 overlays.sort_by_key(|overlay| overlay.committed_version);
221}
222
223pub fn verify_overlays_newest_last(overlays: &[DataOverlayFile]) -> Result<()> {
231 for pair in overlays.windows(2) {
232 if pair[0].committed_version > pair[1].committed_version {
233 return Err(Error::invalid_input(format!(
234 "overlay files must be stored newest-last, but committed_version {} precedes {}",
235 pair[0].committed_version, pair[1].committed_version
236 )));
237 }
238 }
239 Ok(())
240}
241
242pub fn tombstone_overlay_fields(overlays: &mut Vec<DataOverlayFile>, fields: &[u32]) {
253 for overlay in overlays.iter_mut() {
254 let tombstoned: Vec<i32> = overlay
255 .data_file
256 .fields
257 .iter()
258 .map(|&field| {
259 if field >= 0 && fields.contains(&(field as u32)) {
260 TOMBSTONE_FIELD_ID
261 } else {
262 field
263 }
264 })
265 .collect();
266 overlay.data_file.fields = tombstoned.into();
267 }
268 overlays.retain(|overlay| {
269 overlay
270 .data_file
271 .fields
272 .iter()
273 .any(|&field| field != TOMBSTONE_FIELD_ID)
274 });
275}
276
277impl From<&DataOverlayFile> for pb::DataOverlayFile {
278 fn from(overlay: &DataOverlayFile) -> Self {
279 let coverage = match &overlay.coverage {
280 OverlayCoverage::Shared(bitmap) => {
281 pb::data_overlay_file::Coverage::SharedOffsetBitmap(serialize_roaring(bitmap))
282 }
283 OverlayCoverage::PerField(bitmaps) => {
284 pb::data_overlay_file::Coverage::FieldCoverage(pb::FieldCoverage {
285 offset_bitmaps: bitmaps.iter().map(|b| serialize_roaring(b)).collect(),
286 })
287 }
288 };
289 Self {
290 data_file: Some(pb::DataFile::from(&overlay.data_file)),
291 coverage: Some(coverage),
292 committed_version: overlay.committed_version,
293 }
294 }
295}
296
297impl TryFrom<pb::DataOverlayFile> for DataOverlayFile {
298 type Error = Error;
299
300 fn try_from(proto: pb::DataOverlayFile) -> Result<Self> {
301 let data_file = proto
302 .data_file
303 .ok_or_else(|| Error::invalid_input("DataOverlayFile is missing its data_file"))?;
304 let path = Path::from(data_file.path.as_str());
305 let coverage = match proto.coverage {
306 Some(pb::data_overlay_file::Coverage::SharedOffsetBitmap(bytes)) => {
307 OverlayCoverage::Shared(Arc::new(deserialize_roaring(&bytes, &path)?))
308 }
309 Some(pb::data_overlay_file::Coverage::FieldCoverage(fc)) => OverlayCoverage::PerField(
310 fc.offset_bitmaps
311 .iter()
312 .map(|b| deserialize_roaring(b, &path).map(Arc::new))
313 .collect::<Result<_>>()?,
314 ),
315 None => {
316 return Err(Error::invalid_input(
317 "DataOverlayFile is missing its coverage",
318 ));
319 }
320 };
321 Ok(Self {
322 data_file: DataFile::try_from(data_file)?,
323 coverage,
324 committed_version: proto.committed_version,
325 })
326 }
327}
328
329#[cfg(test)]
330mod tests {
331 use super::*;
332
333 #[test]
334 fn test_data_overlay_missing_fields_error() {
335 let no_coverage = pb::DataOverlayFile {
337 data_file: Some(pb::DataFile::from(&DataFile::new_legacy_from_fields(
338 "overlay.lance",
339 vec![3],
340 None,
341 ))),
342 coverage: None,
343 committed_version: 1,
344 };
345 let err = DataOverlayFile::try_from(no_coverage).unwrap_err();
346 assert!(err.to_string().contains("missing its coverage"), "{err}");
347
348 let no_data_file = pb::DataOverlayFile {
349 data_file: None,
350 coverage: Some(pb::data_overlay_file::Coverage::SharedOffsetBitmap(
351 serialize_roaring(&RoaringBitmap::from_iter([0u32])),
352 )),
353 committed_version: 1,
354 };
355 let err = DataOverlayFile::try_from(no_data_file).unwrap_err();
356 assert!(err.to_string().contains("missing its data_file"), "{err}");
357 }
358
359 #[test]
360 fn test_coverage_bitmap_serialized_run_optimized() {
361 let bitmap = RoaringBitmap::from_sorted_iter(0..1_000_000).unwrap();
362 let unoptimized_size = bitmap.serialized_size();
363
364 let overlay = DataOverlayFile {
365 data_file: DataFile::new_legacy_from_fields("overlay.lance", vec![3], None),
366 coverage: OverlayCoverage::dense(bitmap.clone()),
367 committed_version: 1,
368 };
369
370 let proto = pb::DataOverlayFile::from(&overlay);
371 let Some(pb::data_overlay_file::Coverage::SharedOffsetBitmap(bytes)) = &proto.coverage
372 else {
373 panic!("dense coverage must serialize as a shared offset bitmap");
374 };
375 assert!(
376 bytes.len() < unoptimized_size / 100,
377 "expected run-optimized coverage ({} bytes) to be <1% of the \
378 unoptimized serialization ({} bytes)",
379 bytes.len(),
380 unoptimized_size
381 );
382
383 let recovered = DataOverlayFile::try_from(proto).unwrap();
384 assert_eq!(recovered.coverage, OverlayCoverage::dense(bitmap));
385 }
386
387 #[test]
388 fn test_overlay_coverage_serde_json_roundtrip() {
389 for coverage in [
392 OverlayCoverage::dense(RoaringBitmap::from_iter([1u32, 5, 100])),
393 OverlayCoverage::dense(RoaringBitmap::new()),
394 OverlayCoverage::sparse(vec![
395 RoaringBitmap::from_iter([2u32, 3]),
396 RoaringBitmap::new(),
397 ]),
398 OverlayCoverage::sparse(vec![]),
399 ] {
400 let json = serde_json::to_string(&coverage).unwrap();
401 let back: OverlayCoverage = serde_json::from_str(&json).unwrap();
402 assert_eq!(back, coverage);
403 }
404 }
405
406 #[test]
407 fn test_tombstone_overlay_fields() {
408 let mut overlays = vec![
412 DataOverlayFile {
413 data_file: DataFile::new_legacy_from_fields("a.lance", vec![3, 5], None),
414 coverage: OverlayCoverage::sparse(vec![
415 RoaringBitmap::from_iter([0u32]),
416 RoaringBitmap::from_iter([1u32]),
417 ]),
418 committed_version: 1,
419 },
420 DataOverlayFile {
421 data_file: DataFile::new_legacy_from_fields("b.lance", vec![5], None),
422 coverage: OverlayCoverage::dense(RoaringBitmap::from_iter([0u32])),
423 committed_version: 1,
424 },
425 DataOverlayFile {
426 data_file: DataFile::new_legacy_from_fields("c.lance", vec![7], None),
427 coverage: OverlayCoverage::dense(RoaringBitmap::from_iter([0u32])),
428 committed_version: 1,
429 },
430 ];
431
432 tombstone_overlay_fields(&mut overlays, &[5]);
433
434 assert_eq!(overlays.len(), 2);
436 assert_eq!(
438 overlays[0].data_file.fields.as_ref(),
439 &[3, TOMBSTONE_FIELD_ID]
440 );
441 assert_eq!(overlays[1].data_file.fields.as_ref(), &[7]);
443 }
444
445 #[test]
446 fn test_verify_overlays_newest_last() {
447 let mk = |version: u64| DataOverlayFile {
448 data_file: DataFile::new_legacy_from_fields("o.lance", vec![3], None),
449 coverage: OverlayCoverage::dense(RoaringBitmap::from_iter([0u32])),
450 committed_version: version,
451 };
452 assert!(verify_overlays_newest_last(&[]).is_ok());
454 assert!(verify_overlays_newest_last(&[mk(1), mk(2), mk(2), mk(5)]).is_ok());
455 let err = verify_overlays_newest_last(&[mk(2), mk(1)]).unwrap_err();
457 assert!(err.to_string().contains("newest-last"), "{err}");
458 }
459
460 #[test]
461 fn test_coverage_for_field_out_of_bounds() {
462 let overlay = DataOverlayFile {
463 data_file: DataFile::new_legacy_from_fields("o.lance", vec![2, 4], None),
464 coverage: OverlayCoverage::sparse(vec![
465 RoaringBitmap::from_iter([1u32]),
466 RoaringBitmap::from_iter([2u32]),
467 ]),
468 committed_version: 1,
469 };
470 assert!(overlay.coverage_for_field(0).is_ok());
471 assert!(overlay.coverage_for_field(1).is_ok());
472 let err = overlay.coverage_for_field(5).unwrap_err();
473 assert!(err.to_string().contains("field position"), "{err}");
474 }
475}