1use std::{collections::HashSet, sync::Arc};
5
6use arrow_array::{RecordBatch, UInt32Array};
7use arrow_ipc::CompressionType;
8use arrow_ipc::reader::FileReader as ArrowFileReader;
9use arrow_ipc::writer::{FileWriter as ArrowFileWriter, IpcWriteOptions};
10use arrow_schema::{ArrowError, DataType, Field, Schema};
11use bytes::Buf;
12use lance_core::error::{CorruptFileSnafu, box_error};
13use lance_core::utils::deletion::DeletionVector;
14use lance_core::utils::tracing::{AUDIT_MODE_CREATE, AUDIT_TYPE_DELETION, TRACE_FILE_AUDIT};
15use lance_core::{Error, Result};
16use lance_io::object_store::ObjectStore;
17use object_store::path::Path;
18use rand::Rng;
19use roaring::bitmap::RoaringBitmap;
20use snafu::ResultExt;
21use tracing::{info, instrument};
22
23use crate::format::{DeletionFile, DeletionFileType};
24
25pub const DELETIONS_DIR: &str = "_deletions";
26
27fn deletion_arrow_schema() -> Arc<Schema> {
29 Arc::new(Schema::new(vec![Field::new(
30 "row_id",
31 DataType::UInt32,
32 false,
33 )]))
34}
35
36pub fn deletion_file_path(base: &Path, fragment_id: u64, deletion_file: &DeletionFile) -> Path {
38 let DeletionFile {
39 read_version,
40 id,
41 file_type,
42 ..
43 } = deletion_file;
44 let suffix = file_type.suffix();
45 base.child(DELETIONS_DIR)
46 .child(format!("{fragment_id}-{read_version}-{id}.{suffix}"))
47}
48
49pub fn relative_deletion_file_path(fragment_id: u64, deletion_file: &DeletionFile) -> String {
50 let DeletionFile {
51 read_version,
52 id,
53 file_type,
54 ..
55 } = deletion_file;
56 let suffix = file_type.suffix();
57 format!("{DELETIONS_DIR}/{fragment_id}-{read_version}-{id}.{suffix}")
58}
59
60pub async fn write_deletion_file(
65 base: &Path,
66 fragment_id: u64,
67 read_version: u64,
68 removed_rows: &DeletionVector,
69 object_store: &ObjectStore,
70) -> Result<Option<DeletionFile>> {
71 let deletion_file = match removed_rows {
72 DeletionVector::NoDeletions => None,
73 DeletionVector::Set(set) => {
74 let id = rand::rng().random::<u64>();
75 let deletion_file = DeletionFile {
76 read_version,
77 id,
78 file_type: DeletionFileType::Array,
79 num_deleted_rows: Some(set.len()),
80 base_id: None,
81 };
82 let path = deletion_file_path(base, fragment_id, &deletion_file);
83
84 let array = UInt32Array::from_iter(set.iter().copied());
85 let array = Arc::new(array);
86
87 let schema = deletion_arrow_schema();
88 let batch = RecordBatch::try_new(schema.clone(), vec![array])?;
89
90 let mut out: Vec<u8> = Vec::new();
91 let write_options =
92 IpcWriteOptions::default().try_with_compression(Some(CompressionType::ZSTD))?;
93 {
94 let mut writer = ArrowFileWriter::try_new_with_options(
95 &mut out,
96 schema.as_ref(),
97 write_options,
98 )?;
99 writer.write(&batch)?;
100 writer.finish()?;
101 }
103
104 object_store.put(&path, &out).await?;
105
106 info!(target: TRACE_FILE_AUDIT, mode=AUDIT_MODE_CREATE, r#type=AUDIT_TYPE_DELETION, path = path.to_string());
107
108 Some(deletion_file)
109 }
110 DeletionVector::Bitmap(bitmap) => {
111 let id = rand::rng().random::<u64>();
112 let deletion_file = DeletionFile {
113 read_version,
114 id,
115 file_type: DeletionFileType::Bitmap,
116 num_deleted_rows: Some(bitmap.len() as usize),
117 base_id: None,
118 };
119 let path = deletion_file_path(base, fragment_id, &deletion_file);
120
121 let mut out: Vec<u8> = Vec::new();
122 bitmap.serialize_into(&mut out)?;
123
124 object_store.put(&path, &out).await?;
125
126 info!(target: TRACE_FILE_AUDIT, mode=AUDIT_MODE_CREATE, r#type=AUDIT_TYPE_DELETION, path = path.to_string());
127
128 Some(deletion_file)
129 }
130 };
131 Ok(deletion_file)
132}
133
134#[instrument(
135 level = "debug",
136 skip(base, object_store),
137 fields(
138 base = base.as_ref(),
139 bytes_read = tracing::field::Empty
140 )
141)]
142pub async fn read_deletion_file(
143 fragment_id: u64,
144 deletion_file: &DeletionFile,
145 base: &Path,
146 object_store: &ObjectStore,
147) -> Result<DeletionVector> {
148 let span = tracing::Span::current();
149 match deletion_file.file_type {
150 DeletionFileType::Array => {
151 let path = deletion_file_path(base, fragment_id, deletion_file);
152
153 let data = object_store.read_one_all(&path).await?;
154 span.record("bytes_read", data.len());
155 let data = std::io::Cursor::new(data);
156 let mut batches: Vec<RecordBatch> = ArrowFileReader::try_new(data, None)?
157 .collect::<std::result::Result<_, ArrowError>>()
158 .map_err(box_error)
159 .context(CorruptFileSnafu { path: path.clone() })?;
160
161 if batches.len() != 1 {
162 return Err(Error::corrupt_file(
163 path,
164 format!(
165 "Expected exactly one batch in deletion file, got {}",
166 batches.len()
167 ),
168 ));
169 }
170
171 let batch = batches.pop().unwrap();
172 if batch.schema() != deletion_arrow_schema() {
173 return Err(Error::corrupt_file(
174 path,
175 format!(
176 "Expected schema {:?} in deletion file, got {:?}",
177 deletion_arrow_schema(),
178 batch.schema()
179 ),
180 ));
181 }
182
183 let array = batch.columns()[0]
184 .as_any()
185 .downcast_ref::<UInt32Array>()
186 .unwrap();
187
188 let mut set = HashSet::with_capacity(array.len());
189 for val in array.iter() {
190 if let Some(val) = val {
191 set.insert(val);
192 } else {
193 return Err(Error::corrupt_file(
194 path,
195 "Null values are not allowed in deletion files",
196 ));
197 }
198 }
199
200 Ok(DeletionVector::Set(set))
201 }
202 DeletionFileType::Bitmap => {
203 let path = deletion_file_path(base, fragment_id, deletion_file);
204
205 let data = object_store.read_one_all(&path).await?;
206 span.record("bytes_read", data.len());
207 let reader = data.reader();
208 let bitmap = RoaringBitmap::deserialize_from(reader)
209 .map_err(box_error)
210 .context(CorruptFileSnafu { path })?;
211
212 Ok(DeletionVector::Bitmap(bitmap))
213 }
214 }
215}
216
217#[cfg(test)]
218mod test {
219
220 use super::*;
221
222 #[tokio::test]
223 async fn test_write_no_deletions() {
224 let dv = DeletionVector::NoDeletions;
225
226 let (object_store, path) = ObjectStore::from_uri("memory:///no_deletion")
227 .await
228 .unwrap();
229 let file = write_deletion_file(&path, 0, 0, &dv, &object_store)
230 .await
231 .unwrap();
232 assert!(file.is_none());
233 }
234
235 #[tokio::test]
236 async fn test_write_array() {
237 let dv = DeletionVector::Set(HashSet::from_iter(0..100));
238
239 let fragment_id = 21;
240 let read_version = 12;
241
242 let object_store = ObjectStore::memory();
243 let path = Path::from("/write");
244 let file = write_deletion_file(&path, fragment_id, read_version, &dv, &object_store)
245 .await
246 .unwrap();
247
248 assert!(matches!(
249 file,
250 Some(DeletionFile {
251 file_type: DeletionFileType::Array,
252 ..
253 })
254 ));
255
256 let file = file.unwrap();
257 assert_eq!(file.read_version, read_version);
258 let path = deletion_file_path(&path, fragment_id, &file);
259 assert_eq!(
260 path,
261 Path::from(format!("/write/_deletions/21-12-{}.arrow", file.id))
262 );
263
264 let data = object_store
265 .inner
266 .get(&path)
267 .await
268 .unwrap()
269 .bytes()
270 .await
271 .unwrap();
272 let data = std::io::Cursor::new(data);
273 let mut batches: Vec<RecordBatch> = ArrowFileReader::try_new(data, None)
274 .unwrap()
275 .collect::<std::result::Result<_, ArrowError>>()
276 .unwrap();
277
278 assert_eq!(batches.len(), 1);
279 let batch = batches.pop().unwrap();
280 assert_eq!(batch.schema(), deletion_arrow_schema());
281 let array = batch["row_id"]
282 .as_any()
283 .downcast_ref::<UInt32Array>()
284 .unwrap();
285 let read_dv = DeletionVector::from_iter(array.iter().map(|v| v.unwrap()));
286 assert_eq!(read_dv, dv);
287 }
288
289 #[tokio::test]
290 async fn test_write_bitmap() {
291 let dv = DeletionVector::Bitmap(RoaringBitmap::from_iter(0..100));
292
293 let fragment_id = 21;
294 let read_version = 12;
295
296 let object_store = ObjectStore::memory();
297 let path = Path::from("/bitmap");
298 let file = write_deletion_file(&path, fragment_id, read_version, &dv, &object_store)
299 .await
300 .unwrap();
301
302 assert!(matches!(
303 file,
304 Some(DeletionFile {
305 file_type: DeletionFileType::Bitmap,
306 ..
307 })
308 ));
309
310 let file = file.unwrap();
311 assert_eq!(file.read_version, read_version);
312 let path = deletion_file_path(&path, fragment_id, &file);
313 assert_eq!(
314 path,
315 Path::from(format!("/bitmap/_deletions/21-12-{}.bin", file.id))
316 );
317
318 let data = object_store
319 .inner
320 .get(&path)
321 .await
322 .unwrap()
323 .bytes()
324 .await
325 .unwrap();
326 let reader = data.reader();
327 let read_bitmap = RoaringBitmap::deserialize_from(reader).unwrap();
328 assert_eq!(read_bitmap, dv.into_iter().collect::<RoaringBitmap>());
329 }
330
331 #[tokio::test]
332 async fn test_roundtrip_array() {
333 let dv = DeletionVector::Set(HashSet::from_iter(0..100));
334
335 let fragment_id = 21;
336 let read_version = 12;
337
338 let object_store = ObjectStore::memory();
339 let path = Path::from("/roundtrip");
340 let file = write_deletion_file(&path, fragment_id, read_version, &dv, &object_store)
341 .await
342 .unwrap();
343
344 let read_dv = read_deletion_file(fragment_id, &file.unwrap(), &path, &object_store)
345 .await
346 .unwrap();
347 assert_eq!(read_dv, dv);
348 }
349
350 #[tokio::test]
351 async fn test_roundtrip_bitmap() {
352 let dv = DeletionVector::Bitmap(RoaringBitmap::from_iter(0..100));
353
354 let fragment_id = 21;
355 let read_version = 12;
356
357 let object_store = ObjectStore::memory();
358 let path = Path::from("/bitmap");
359 let file = write_deletion_file(&path, fragment_id, read_version, &dv, &object_store)
360 .await
361 .unwrap();
362
363 let read_dv = read_deletion_file(fragment_id, &file.unwrap(), &path, &object_store)
364 .await
365 .unwrap();
366 assert_eq!(read_dv, dv);
367 }
368}