1use ailake_catalog::provider::DeletionVector;
11use ailake_core::{AilakeError, AilakeResult};
12use ailake_store::Store;
13use arrow_array::RecordBatch;
14use roaring::RoaringBitmap;
15use std::sync::Arc;
16
17pub async fn load_deletion_vector(
28 store: &Arc<dyn Store>,
29 dv: &DeletionVector,
30) -> AilakeResult<RoaringBitmap> {
31 let bytes = store
32 .get_range(&dv.path, dv.offset..dv.offset + dv.length)
33 .await?;
34
35 RoaringBitmap::deserialize_from(bytes.as_ref()).map_err(|e| {
36 AilakeError::Io(std::io::Error::other(format!(
37 "ailake: failed to deserialize Deletion Vector bitmap from '{}' \
38 (offset={}, length={}): {e}",
39 dv.path, dv.offset, dv.length
40 )))
41 })
42}
43
44pub fn filter_deleted_rows<T>(
54 batch: RecordBatch,
55 parallel: Vec<T>,
56 bitmap: &RoaringBitmap,
57) -> AilakeResult<(RecordBatch, Vec<T>)> {
58 if bitmap.is_empty() {
59 return Ok((batch, parallel));
60 }
61 let n = batch.num_rows();
62 let keep: Vec<bool> = (0..n).map(|i| !bitmap.contains(i as u32)).collect();
63 let filtered_parallel: Vec<T> = parallel
64 .into_iter()
65 .zip(keep.iter())
66 .filter_map(|(v, &k)| k.then_some(v))
67 .collect();
68 let mask = arrow_array::BooleanArray::from(keep);
69 let filtered_batch = arrow_select::filter::filter_record_batch(&batch, &mask)
70 .map_err(|e| AilakeError::Arrow(e.to_string()))?;
71 Ok((filtered_batch, filtered_parallel))
72}
73
74#[inline]
78pub fn has_deletions(bitmap: &RoaringBitmap, row_ids: &[u64]) -> bool {
79 row_ids.iter().any(|&id| bitmap.contains(id as u32))
80}
81
82#[cfg(test)]
83mod tests {
84 use super::*;
85 use ailake_store::LocalStore;
86 use bytes::Bytes;
87 use roaring::RoaringBitmap;
88
89 fn make_bitmap_bytes(deleted: &[u32]) -> Vec<u8> {
90 let mut bm = RoaringBitmap::new();
91 for &r in deleted {
92 bm.insert(r);
93 }
94 let mut buf = Vec::new();
95 bm.serialize_into(&mut buf).unwrap();
96 buf
97 }
98
99 #[tokio::test]
100 async fn load_dv_roundtrip() {
101 let dir = tempfile::tempdir().unwrap();
102 let bitmap_bytes = make_bitmap_bytes(&[0, 5, 42, 1000]);
103
104 let offset: u64 = 16; let mut file_bytes = vec![0u8; offset as usize]; file_bytes.extend_from_slice(&bitmap_bytes);
110
111 let dvd_path = "data/dv-0001.dvd";
112 let store: Arc<dyn Store> = Arc::new(LocalStore::new(dir.path()));
113 store.put(dvd_path, Bytes::from(file_bytes)).await.unwrap();
114
115 let dv = DeletionVector {
116 path: dvd_path.to_string(),
117 offset,
118 length: bitmap_bytes.len() as u64,
119 cardinality: 4,
120 };
121
122 let bm = load_deletion_vector(&store, &dv).await.unwrap();
123 assert!(bm.contains(0));
124 assert!(bm.contains(5));
125 assert!(bm.contains(42));
126 assert!(bm.contains(1000));
127 assert!(!bm.contains(1)); assert_eq!(bm.len(), 4);
129 }
130
131 #[test]
132 fn has_deletions_detects_overlap() {
133 let mut bm = RoaringBitmap::new();
134 bm.insert(10);
135 bm.insert(20);
136
137 assert!(has_deletions(&bm, &[5, 10, 15])); assert!(!has_deletions(&bm, &[1, 2, 3])); }
140}