Skip to main content

lance_encoding/encodings/logical/
blob.rs

1// SPDX-License-Identifier: Apache-2.0
2// SPDX-FileCopyrightText: Copyright The Lance Authors
3
4use std::{collections::HashMap, sync::Arc};
5
6use arrow_array::{
7    Array, ArrayRef, StructArray, UInt64Array,
8    builder::{PrimitiveBuilder, StringBuilder},
9    cast::AsArray,
10    types::{UInt8Type, UInt32Type, UInt64Type},
11};
12use arrow_buffer::Buffer;
13use arrow_schema::{DataType, Field as ArrowField, Fields};
14use futures::{FutureExt, future::BoxFuture};
15use lance_core::{
16    Error, Result, datatypes::BLOB_V2_DESC_FIELDS, datatypes::Field, error::LanceOptionExt,
17};
18
19use crate::{
20    buffer::LanceBuffer,
21    constants::PACKED_STRUCT_META_KEY,
22    decoder::PageEncoding,
23    encoder::{EncodeTask, EncodedColumn, EncodedPage, FieldEncoder, OutOfLineBuffers},
24    format::ProtobufUtils21,
25    repdef::{DefinitionInterpretation, RepDefBuilder},
26};
27use lance_core::datatypes::BlobKind;
28
29/// Blob structural encoder - stores large binary data in external buffers
30///
31/// This encoder takes large binary arrays and stores them outside the normal
32/// page structure. It creates a descriptor (position, size) for each blob
33/// that is stored inline in the page.
34pub struct BlobStructuralEncoder {
35    // Encoder for the descriptors (position/size struct)
36    descriptor_encoder: Box<dyn FieldEncoder>,
37    // Set when we first see data
38    def_meaning: Option<Arc<[DefinitionInterpretation]>>,
39}
40
41impl BlobStructuralEncoder {
42    pub fn new(
43        field: &Field,
44        make_descriptor_encoder: impl FnOnce(Field) -> Result<Box<dyn FieldEncoder>>,
45    ) -> Result<Self> {
46        // Create descriptor field: struct<position: u64, size: u64>
47        // Preserve the original field's metadata for packed struct
48        let mut descriptor_metadata = HashMap::with_capacity(1);
49        descriptor_metadata.insert(PACKED_STRUCT_META_KEY.to_string(), "true".to_string());
50
51        let descriptor_data_type = DataType::Struct(Fields::from(vec![
52            ArrowField::new("position", DataType::UInt64, false),
53            ArrowField::new("size", DataType::UInt64, false),
54        ]));
55
56        // Use the original field's name for the descriptor
57        let descriptor_field = Field::try_from(
58            ArrowField::new(&field.name, descriptor_data_type, field.nullable)
59                .with_metadata(descriptor_metadata),
60        )?;
61
62        // Use PrimitiveStructuralEncoder to handle the descriptor
63        let descriptor_encoder = make_descriptor_encoder(descriptor_field)?;
64
65        Ok(Self {
66            descriptor_encoder,
67            def_meaning: None,
68        })
69    }
70
71    fn wrap_tasks(
72        tasks: Vec<EncodeTask>,
73        def_meaning: Arc<[DefinitionInterpretation]>,
74    ) -> Vec<EncodeTask> {
75        tasks
76            .into_iter()
77            .map(|task| {
78                let def_meaning = def_meaning.clone();
79                task.then(|encoded_page| async move {
80                    let encoded_page = encoded_page?;
81
82                    let PageEncoding::Structural(inner_layout) = encoded_page.description else {
83                        return Err(Error::internal(
84                            "Expected inner encoding to return structural layout".to_string(),
85                        ));
86                    };
87
88                    let wrapped = ProtobufUtils21::blob_layout(inner_layout, &def_meaning);
89                    Ok(EncodedPage {
90                        column_idx: encoded_page.column_idx,
91                        data: encoded_page.data,
92                        description: PageEncoding::Structural(wrapped),
93                        num_rows: encoded_page.num_rows,
94                        row_number: encoded_page.row_number,
95                    })
96                })
97                .boxed()
98            })
99            .collect::<Vec<_>>()
100    }
101}
102
103impl FieldEncoder for BlobStructuralEncoder {
104    fn maybe_encode(
105        &mut self,
106        array: ArrayRef,
107        external_buffers: &mut OutOfLineBuffers,
108        mut repdef: RepDefBuilder,
109        row_number: u64,
110        num_rows: u64,
111    ) -> Result<Vec<EncodeTask>> {
112        if let Some(validity) = array.nulls() {
113            repdef.add_validity_bitmap(validity.clone());
114        } else {
115            repdef.add_no_null(array.len());
116        }
117
118        // Convert input array to LargeBinary
119        let binary_array = array.as_binary_opt::<i64>().ok_or_else(|| {
120            Error::invalid_input_source(
121                format!("Expected LargeBinary array, got {}", array.data_type()).into(),
122            )
123        })?;
124
125        let repdef = RepDefBuilder::serialize(vec![repdef]);
126
127        let rep = repdef.repetition_levels.as_ref();
128        let def = repdef.definition_levels.as_ref();
129        let def_meaning: Arc<[DefinitionInterpretation]> = repdef.def_meaning.into();
130
131        // A blob page stores one definition interpretation for all of its rows.
132        // The descriptor encoder can buffer multiple input arrays, so finish the
133        // pending page before a later array changes from all-valid to nullable (or
134        // vice versa).
135        let mut encode_tasks = match self.def_meaning.as_ref() {
136            Some(existing) if existing != &def_meaning => {
137                let existing = existing.clone();
138                Self::wrap_tasks(self.descriptor_encoder.flush(external_buffers)?, existing)
139            }
140            _ => Vec::new(),
141        };
142        self.def_meaning = Some(def_meaning.clone());
143
144        // Collect positions and sizes
145        let mut positions = Vec::with_capacity(binary_array.len());
146        let mut sizes = Vec::with_capacity(binary_array.len());
147
148        for i in 0..binary_array.len() {
149            if binary_array.is_null(i) {
150                // Null values are smuggled into the positions array
151
152                // If we have null values we must have definition levels
153                let mut repdef = (def.expect_ok()?[i] as u64) << 16;
154                if let Some(rep) = rep {
155                    repdef += rep[i] as u64;
156                }
157
158                debug_assert_ne!(repdef, 0);
159                positions.push(repdef);
160                sizes.push(0);
161            } else {
162                let value = binary_array.value(i);
163                if value.is_empty() {
164                    // Empty values
165                    positions.push(0);
166                    sizes.push(0);
167                } else {
168                    // Add data to external buffers
169                    let position =
170                        external_buffers.add_buffer(LanceBuffer::from(Buffer::from(value)));
171                    positions.push(position);
172                    sizes.push(value.len() as u64);
173                }
174            }
175        }
176
177        // Create descriptor array
178        let position_array = Arc::new(UInt64Array::from(positions));
179        let size_array = Arc::new(UInt64Array::from(sizes));
180        let descriptor_array = Arc::new(StructArray::new(
181            Fields::from(vec![
182                ArrowField::new("position", DataType::UInt64, false),
183                ArrowField::new("size", DataType::UInt64, false),
184            ]),
185            vec![position_array as ArrayRef, size_array as ArrayRef],
186            None, // Descriptors are never null
187        ));
188
189        // Delegate to descriptor encoder
190        let descriptor_tasks = self.descriptor_encoder.maybe_encode(
191            descriptor_array,
192            external_buffers,
193            RepDefBuilder::default(),
194            row_number,
195            num_rows,
196        )?;
197        encode_tasks.extend(Self::wrap_tasks(descriptor_tasks, def_meaning));
198
199        Ok(encode_tasks)
200    }
201
202    fn flush(&mut self, external_buffers: &mut OutOfLineBuffers) -> Result<Vec<EncodeTask>> {
203        let encode_tasks = self.descriptor_encoder.flush(external_buffers)?;
204
205        // Use the cached def meaning.  If we haven't seen any data yet then we can just use a dummy
206        // value (not clear there would be any encode tasks in that case)
207        let def_meaning = self
208            .def_meaning
209            .clone()
210            .unwrap_or_else(|| Arc::new([DefinitionInterpretation::AllValidItem]));
211
212        Ok(Self::wrap_tasks(encode_tasks, def_meaning))
213    }
214
215    fn finish(
216        &mut self,
217        external_buffers: &mut OutOfLineBuffers,
218    ) -> BoxFuture<'_, Result<Vec<EncodedColumn>>> {
219        self.descriptor_encoder.finish(external_buffers)
220    }
221
222    fn num_columns(&self) -> u32 {
223        self.descriptor_encoder.num_columns()
224    }
225}
226
227/// Blob v2 structural encoder
228pub struct BlobV2StructuralEncoder {
229    descriptor_encoder: Box<dyn FieldEncoder>,
230}
231
232impl BlobV2StructuralEncoder {
233    pub fn new(
234        field: &Field,
235        make_descriptor_encoder: impl FnOnce(Field) -> Result<Box<dyn FieldEncoder>>,
236    ) -> Result<Self> {
237        let mut descriptor_metadata = HashMap::with_capacity(1);
238        descriptor_metadata.insert(PACKED_STRUCT_META_KEY.to_string(), "true".to_string());
239
240        let descriptor_data_type = DataType::Struct(BLOB_V2_DESC_FIELDS.clone());
241
242        let descriptor_field = Field::try_from(
243            ArrowField::new(&field.name, descriptor_data_type, field.nullable)
244                .with_metadata(descriptor_metadata),
245        )?;
246
247        let descriptor_encoder = make_descriptor_encoder(descriptor_field)?;
248
249        Ok(Self { descriptor_encoder })
250    }
251}
252
253impl FieldEncoder for BlobV2StructuralEncoder {
254    fn maybe_encode(
255        &mut self,
256        array: ArrayRef,
257        external_buffers: &mut OutOfLineBuffers,
258        repdef: RepDefBuilder,
259        row_number: u64,
260        num_rows: u64,
261    ) -> Result<Vec<EncodeTask>> {
262        let struct_arr = array.as_struct();
263
264        let kind_col = struct_arr
265            .column_by_name("kind")
266            .ok_or_else(|| {
267                Error::invalid_input_source("Blob v2 struct missing `kind` field".into())
268            })?
269            .as_primitive::<UInt8Type>();
270        let data_col = struct_arr
271            .column_by_name("data")
272            .ok_or_else(|| {
273                Error::invalid_input_source("Blob v2 struct missing `data` field".into())
274            })?
275            .as_binary::<i64>();
276        let uri_col = struct_arr
277            .column_by_name("uri")
278            .ok_or_else(|| {
279                Error::invalid_input_source("Blob v2 struct missing `uri` field".into())
280            })?
281            .as_string::<i32>();
282        let blob_id_col = struct_arr
283            .column_by_name("blob_id")
284            .ok_or_else(|| {
285                Error::invalid_input_source("Blob v2 struct missing `blob_id` field".into())
286            })?
287            .as_primitive::<UInt32Type>();
288        let blob_size_col = struct_arr
289            .column_by_name("blob_size")
290            .ok_or_else(|| {
291                Error::invalid_input_source("Blob v2 struct missing `blob_size` field".into())
292            })?
293            .as_primitive::<UInt64Type>();
294        let packed_position_col = struct_arr
295            .column_by_name("position")
296            .ok_or_else(|| {
297                Error::invalid_input_source("Blob v2 struct missing `position` field".into())
298            })?
299            .as_primitive::<UInt64Type>();
300
301        let row_count = struct_arr.len();
302
303        let mut kind_builder = PrimitiveBuilder::<UInt8Type>::with_capacity(row_count);
304        let mut position_builder = PrimitiveBuilder::<UInt64Type>::with_capacity(row_count);
305        let mut size_builder = PrimitiveBuilder::<UInt64Type>::with_capacity(row_count);
306        let mut blob_id_builder = PrimitiveBuilder::<UInt32Type>::with_capacity(row_count);
307        let mut uri_builder = StringBuilder::with_capacity(row_count, row_count * 16);
308
309        for i in 0..row_count {
310            let (kind_value, position_value, size_value, blob_id_value, uri_value) =
311                if struct_arr.is_null(i) || kind_col.is_null(i) {
312                    (BlobKind::Inline as u8, 0, 0, 0, "".to_string())
313                } else {
314                    let kind_val = BlobKind::try_from(kind_col.value(i))?;
315                    match kind_val {
316                        BlobKind::Dedicated => (
317                            BlobKind::Dedicated as u8,
318                            0,
319                            blob_size_col.value(i),
320                            blob_id_col.value(i),
321                            "".to_string(),
322                        ),
323                        BlobKind::External => {
324                            let uri = uri_col.value(i).to_string();
325                            let position = if packed_position_col.is_null(i) {
326                                0
327                            } else {
328                                packed_position_col.value(i)
329                            };
330                            let size = if blob_size_col.is_null(i) {
331                                0
332                            } else {
333                                blob_size_col.value(i)
334                            };
335                            let external_base_id = if blob_id_col.is_null(i) {
336                                0
337                            } else {
338                                blob_id_col.value(i)
339                            };
340                            (
341                                BlobKind::External as u8,
342                                position,
343                                size,
344                                external_base_id,
345                                uri,
346                            )
347                        }
348                        BlobKind::Packed => (
349                            BlobKind::Packed as u8,
350                            packed_position_col.value(i),
351                            blob_size_col.value(i),
352                            blob_id_col.value(i),
353                            "".to_string(),
354                        ),
355                        BlobKind::Inline => {
356                            let data_val = data_col.value(i);
357                            let blob_len = data_val.len() as u64;
358                            let position = external_buffers
359                                .add_buffer(LanceBuffer::from(Buffer::from(data_val)));
360
361                            (
362                                BlobKind::Inline as u8,
363                                position,
364                                blob_len,
365                                0,
366                                "".to_string(),
367                            )
368                        }
369                    }
370                };
371
372            kind_builder.append_value(kind_value);
373            position_builder.append_value(position_value);
374            size_builder.append_value(size_value);
375            blob_id_builder.append_value(blob_id_value);
376            uri_builder.append_value(uri_value);
377        }
378        let children: Vec<ArrayRef> = vec![
379            Arc::new(kind_builder.finish()),
380            Arc::new(position_builder.finish()),
381            Arc::new(size_builder.finish()),
382            Arc::new(blob_id_builder.finish()),
383            Arc::new(uri_builder.finish()),
384        ];
385
386        let descriptor_array = Arc::new(StructArray::try_new(
387            BLOB_V2_DESC_FIELDS.clone(),
388            children,
389            struct_arr.nulls().cloned(),
390        )?) as ArrayRef;
391
392        self.descriptor_encoder.maybe_encode(
393            descriptor_array,
394            external_buffers,
395            repdef,
396            row_number,
397            num_rows,
398        )
399    }
400
401    fn flush(&mut self, external_buffers: &mut OutOfLineBuffers) -> Result<Vec<EncodeTask>> {
402        self.descriptor_encoder.flush(external_buffers)
403    }
404
405    fn finish(
406        &mut self,
407        external_buffers: &mut OutOfLineBuffers,
408    ) -> BoxFuture<'_, Result<Vec<EncodedColumn>>> {
409        self.descriptor_encoder.finish(external_buffers)
410    }
411
412    fn num_columns(&self) -> u32 {
413        self.descriptor_encoder.num_columns()
414    }
415}
416
417#[cfg(test)]
418mod tests {
419    use super::*;
420    use crate::{
421        encoder::{ColumnIndexSequence, EncodingOptions},
422        testing::{
423            TestCases, TestEncoding, check_round_trip_encoding_of_data,
424            check_round_trip_encoding_of_data_with_expected, create_test_field_encoder,
425            test_encoding_strategy,
426        },
427    };
428    use arrow_array::{
429        ArrayRef, LargeBinaryArray, StringArray, StructArray, UInt8Array, UInt32Array, UInt64Array,
430    };
431    use arrow_schema::{DataType, Field as ArrowField};
432
433    #[test]
434    fn test_blob_encoder_creation() {
435        let field = Field::try_from(
436            ArrowField::new("blob_field", DataType::LargeBinary, true).with_metadata(
437                HashMap::from([(lance_arrow::BLOB_META_KEY.to_string(), "true".to_string())]),
438            ),
439        )
440        .unwrap();
441        let mut column_index = ColumnIndexSequence::default();
442        let options = EncodingOptions::default();
443        let strategy = test_encoding_strategy(TestEncoding::StructuralU16);
444
445        let encoder =
446            create_test_field_encoder(strategy.as_ref(), &field, &mut column_index, &options);
447
448        assert!(encoder.is_ok());
449    }
450
451    #[tokio::test]
452    async fn test_blob_encoding_simple() {
453        let field = Field::try_from(
454            ArrowField::new("blob_field", DataType::LargeBinary, true).with_metadata(
455                HashMap::from([(lance_arrow::BLOB_META_KEY.to_string(), "true".to_string())]),
456            ),
457        )
458        .unwrap();
459        let mut column_index = ColumnIndexSequence::default();
460        let options = EncodingOptions::default();
461        let strategy = test_encoding_strategy(TestEncoding::StructuralU16);
462
463        let mut encoder =
464            create_test_field_encoder(strategy.as_ref(), &field, &mut column_index, &options)
465                .unwrap();
466
467        // Create test data with larger blobs
468        let large_data = vec![0u8; 1024 * 100]; // 100KB blob
469        let data: Vec<Option<&[u8]>> =
470            vec![Some(b"hello world"), None, Some(&large_data), Some(b"")];
471        let array = Arc::new(LargeBinaryArray::from(data));
472
473        // Test encoding
474        let mut external_buffers = OutOfLineBuffers::new(0, 8);
475        let repdef = RepDefBuilder::default();
476
477        let tasks = encoder
478            .maybe_encode(array, &mut external_buffers, repdef, 0, 4)
479            .unwrap();
480
481        // If no tasks yet, flush to force encoding
482        if tasks.is_empty() {
483            let _flush_tasks = encoder.flush(&mut external_buffers).unwrap();
484        }
485
486        // Should produce encode tasks for the descriptor (or we need more data)
487        // For now, just verify no errors occurred
488        assert!(encoder.num_columns() > 0);
489
490        // Verify external buffers were used for large data
491        let buffers = external_buffers.take_buffers();
492        assert!(
493            !buffers.is_empty(),
494            "Large blobs should be stored in external buffers"
495        );
496    }
497
498    #[tokio::test]
499    async fn test_blob_round_trip() {
500        // Test round-trip encoding with blob metadata
501        let blob_metadata =
502            HashMap::from([(lance_arrow::BLOB_META_KEY.to_string(), "true".to_string())]);
503
504        // Create test data
505        let val1: &[u8] = &vec![1u8; 1024]; // 1KB
506        let val2: &[u8] = &vec![2u8; 10240]; // 10KB
507        let val3: &[u8] = &vec![3u8; 102400]; // 100KB
508        let array = Arc::new(LargeBinaryArray::from(vec![
509            Some(val1),
510            None,
511            Some(val2),
512            Some(val3),
513        ]));
514
515        // Use the standard test harness
516        check_round_trip_encoding_of_data(
517            vec![array],
518            &TestCases::default().with_array_and_u16_encodings(),
519            blob_metadata,
520        )
521        .await;
522    }
523
524    #[tokio::test]
525    async fn test_blob_round_trip_empty_values() {
526        // Empty values share size == 0 with nulls in the descriptor layout
527        // and schedule no read; each must decode to zero-length bytes without
528        // consuming the read result of a following non-empty blob. Empties
529        // are placed before payloads so a misassignment corrupts the output
530        // instead of only exhausting the read iterator.
531        let blob_metadata =
532            HashMap::from([(lance_arrow::BLOB_META_KEY.to_string(), "true".to_string())]);
533
534        let val1: &[u8] = &vec![1u8; 1024];
535        let val2: &[u8] = &vec![2u8; 10240];
536        let empty: &[u8] = &[];
537        let array = Arc::new(LargeBinaryArray::from(vec![
538            Some(empty),
539            Some(val1),
540            None,
541            Some(empty),
542            Some(val2),
543            None,
544            Some(empty),
545        ]));
546
547        check_round_trip_encoding_of_data(vec![array], &TestCases::default(), blob_metadata).await;
548    }
549
550    #[tokio::test]
551    async fn test_blob_round_trip_varying_chunk_nullability() {
552        let blob_metadata =
553            HashMap::from([(lance_arrow::BLOB_META_KEY.to_string(), "true".to_string())]);
554        let all_valid = Arc::new(LargeBinaryArray::from(vec![Some(b"first".as_ref())]));
555        let with_null = Arc::new(LargeBinaryArray::from(vec![
556            Some(b"second".as_ref()),
557            None,
558            Some(b"".as_ref()),
559        ]));
560        let all_valid_again = Arc::new(LargeBinaryArray::from(vec![Some(b"last".as_ref())]));
561
562        check_round_trip_encoding_of_data(
563            vec![all_valid, with_null, all_valid_again],
564            &TestCases::default().with_encoding(TestEncoding::StructuralU16),
565            blob_metadata,
566        )
567        .await;
568    }
569
570    #[tokio::test]
571    async fn test_blob_v2_external_round_trip() {
572        let blob_metadata = HashMap::from([(
573            lance_arrow::ARROW_EXT_NAME_KEY.to_string(),
574            lance_arrow::BLOB_V2_EXT_NAME.to_string(),
575        )]);
576
577        let kind_field = Arc::new(ArrowField::new("kind", DataType::UInt8, true));
578        let data_field = Arc::new(ArrowField::new("data", DataType::LargeBinary, true));
579        let uri_field = Arc::new(ArrowField::new("uri", DataType::Utf8, true));
580        let blob_id_field = Arc::new(ArrowField::new("blob_id", DataType::UInt32, true));
581        let blob_size_field = Arc::new(ArrowField::new("blob_size", DataType::UInt64, true));
582        let position_field = Arc::new(ArrowField::new("position", DataType::UInt64, true));
583
584        let kind_array = UInt8Array::from(vec![
585            BlobKind::Inline as u8,
586            BlobKind::External as u8,
587            BlobKind::External as u8,
588        ]);
589        let data_array = LargeBinaryArray::from(vec![Some(b"inline".as_ref()), None, None]);
590        let uri_array = StringArray::from(vec![
591            None,
592            Some("file:///tmp/external.bin"),
593            Some("s3://bucket/blob"),
594        ]);
595        let blob_id_array = UInt32Array::from(vec![0, 0, 0]);
596        let blob_size_array = UInt64Array::from(vec![0, 0, 0]);
597        let position_array = UInt64Array::from(vec![0, 0, 0]);
598
599        let struct_array = StructArray::from(vec![
600            (kind_field, Arc::new(kind_array) as ArrayRef),
601            (data_field, Arc::new(data_array) as ArrayRef),
602            (uri_field, Arc::new(uri_array) as ArrayRef),
603            (blob_id_field, Arc::new(blob_id_array) as ArrayRef),
604            (blob_size_field, Arc::new(blob_size_array) as ArrayRef),
605            (position_field, Arc::new(position_array) as ArrayRef),
606        ]);
607
608        let expected_descriptor = StructArray::from(vec![
609            (
610                Arc::new(ArrowField::new("kind", DataType::UInt8, false)),
611                Arc::new(UInt8Array::from(vec![
612                    BlobKind::Inline as u8,
613                    BlobKind::External as u8,
614                    BlobKind::External as u8,
615                ])) as ArrayRef,
616            ),
617            (
618                Arc::new(ArrowField::new("position", DataType::UInt64, false)),
619                Arc::new(UInt64Array::from(vec![0, 0, 0])) as ArrayRef,
620            ),
621            (
622                Arc::new(ArrowField::new("size", DataType::UInt64, false)),
623                Arc::new(UInt64Array::from(vec![6, 0, 0])) as ArrayRef,
624            ),
625            (
626                Arc::new(ArrowField::new("blob_id", DataType::UInt32, false)),
627                Arc::new(UInt32Array::from(vec![0, 0, 0])) as ArrayRef,
628            ),
629            (
630                Arc::new(ArrowField::new("blob_uri", DataType::Utf8, false)),
631                Arc::new(StringArray::from(vec![
632                    "",
633                    "file:///tmp/external.bin",
634                    "s3://bucket/blob",
635                ])) as ArrayRef,
636            ),
637        ]);
638
639        check_round_trip_encoding_of_data_with_expected(
640            vec![Arc::new(struct_array)],
641            Some(Arc::new(expected_descriptor)),
642            &TestCases::default().with_u32_structural_encodings(),
643            blob_metadata,
644        )
645        .await;
646    }
647
648    #[tokio::test]
649    async fn test_blob_v2_dedicated_round_trip() {
650        let blob_metadata = HashMap::from([(
651            lance_arrow::ARROW_EXT_NAME_KEY.to_string(),
652            lance_arrow::BLOB_V2_EXT_NAME.to_string(),
653        )]);
654
655        let kind_field = Arc::new(ArrowField::new("kind", DataType::UInt8, true));
656        let data_field = Arc::new(ArrowField::new("data", DataType::LargeBinary, true));
657        let uri_field = Arc::new(ArrowField::new("uri", DataType::Utf8, true));
658        let blob_id_field = Arc::new(ArrowField::new("blob_id", DataType::UInt32, true));
659        let blob_size_field = Arc::new(ArrowField::new("blob_size", DataType::UInt64, true));
660        let position_field = Arc::new(ArrowField::new("position", DataType::UInt64, true));
661
662        let kind_array = UInt8Array::from(vec![BlobKind::Dedicated as u8, BlobKind::Inline as u8]);
663        let data_array = LargeBinaryArray::from(vec![None, Some(b"abc".as_ref())]);
664        let uri_array = StringArray::from(vec![Option::<&str>::None, None]);
665        let blob_id_array = UInt32Array::from(vec![42, 0]);
666        let blob_size_array = UInt64Array::from(vec![12, 0]);
667        let position_array = UInt64Array::from(vec![0, 0]);
668
669        let struct_array = StructArray::from(vec![
670            (kind_field, Arc::new(kind_array) as ArrayRef),
671            (data_field, Arc::new(data_array) as ArrayRef),
672            (uri_field, Arc::new(uri_array) as ArrayRef),
673            (blob_id_field, Arc::new(blob_id_array) as ArrayRef),
674            (blob_size_field, Arc::new(blob_size_array) as ArrayRef),
675            (position_field, Arc::new(position_array) as ArrayRef),
676        ]);
677
678        let expected_descriptor = StructArray::from(vec![
679            (
680                Arc::new(ArrowField::new("kind", DataType::UInt8, false)),
681                Arc::new(UInt8Array::from(vec![
682                    BlobKind::Dedicated as u8,
683                    BlobKind::Inline as u8,
684                ])) as ArrayRef,
685            ),
686            (
687                Arc::new(ArrowField::new("position", DataType::UInt64, false)),
688                Arc::new(UInt64Array::from(vec![0, 0])) as ArrayRef,
689            ),
690            (
691                Arc::new(ArrowField::new("size", DataType::UInt64, false)),
692                Arc::new(UInt64Array::from(vec![12, 3])) as ArrayRef,
693            ),
694            (
695                Arc::new(ArrowField::new("blob_id", DataType::UInt32, false)),
696                Arc::new(UInt32Array::from(vec![42, 0])) as ArrayRef,
697            ),
698            (
699                Arc::new(ArrowField::new("blob_uri", DataType::Utf8, false)),
700                Arc::new(StringArray::from(vec!["", ""])) as ArrayRef,
701            ),
702        ]);
703
704        check_round_trip_encoding_of_data_with_expected(
705            vec![Arc::new(struct_array)],
706            Some(Arc::new(expected_descriptor)),
707            &TestCases::default().with_u32_structural_encodings(),
708            blob_metadata,
709        )
710        .await;
711    }
712
713    #[tokio::test]
714    async fn test_blob_v2_external_with_range_round_trip() {
715        let blob_metadata = HashMap::from([(
716            lance_arrow::ARROW_EXT_NAME_KEY.to_string(),
717            lance_arrow::BLOB_V2_EXT_NAME.to_string(),
718        )]);
719
720        let kind_field = Arc::new(ArrowField::new("kind", DataType::UInt8, true));
721        let data_field = Arc::new(ArrowField::new("data", DataType::LargeBinary, true));
722        let uri_field = Arc::new(ArrowField::new("uri", DataType::Utf8, true));
723        let blob_id_field = Arc::new(ArrowField::new("blob_id", DataType::UInt32, true));
724        let blob_size_field = Arc::new(ArrowField::new("blob_size", DataType::UInt64, true));
725        let position_field = Arc::new(ArrowField::new("position", DataType::UInt64, true));
726
727        let kind_array = UInt8Array::from(vec![BlobKind::External as u8]);
728        let data_array = LargeBinaryArray::from(vec![None::<&[u8]>]);
729        let uri_array = StringArray::from(vec![Some("memory://container.pack")]);
730        let blob_id_array = UInt32Array::from(vec![0]);
731        let blob_size_array = UInt64Array::from(vec![42]);
732        let position_array = UInt64Array::from(vec![7]);
733
734        let struct_array = StructArray::from(vec![
735            (kind_field, Arc::new(kind_array) as ArrayRef),
736            (data_field, Arc::new(data_array) as ArrayRef),
737            (uri_field, Arc::new(uri_array) as ArrayRef),
738            (blob_id_field, Arc::new(blob_id_array) as ArrayRef),
739            (blob_size_field, Arc::new(blob_size_array) as ArrayRef),
740            (position_field, Arc::new(position_array) as ArrayRef),
741        ]);
742
743        let expected_descriptor = StructArray::from(vec![
744            (
745                Arc::new(ArrowField::new("kind", DataType::UInt8, false)),
746                Arc::new(UInt8Array::from(vec![BlobKind::External as u8])) as ArrayRef,
747            ),
748            (
749                Arc::new(ArrowField::new("position", DataType::UInt64, false)),
750                Arc::new(UInt64Array::from(vec![7])) as ArrayRef,
751            ),
752            (
753                Arc::new(ArrowField::new("size", DataType::UInt64, false)),
754                Arc::new(UInt64Array::from(vec![42])) as ArrayRef,
755            ),
756            (
757                Arc::new(ArrowField::new("blob_id", DataType::UInt32, false)),
758                Arc::new(UInt32Array::from(vec![0])) as ArrayRef,
759            ),
760            (
761                Arc::new(ArrowField::new("blob_uri", DataType::Utf8, false)),
762                Arc::new(StringArray::from(vec!["memory://container.pack"])) as ArrayRef,
763            ),
764        ]);
765
766        check_round_trip_encoding_of_data_with_expected(
767            vec![Arc::new(struct_array)],
768            Some(Arc::new(expected_descriptor)),
769            &TestCases::default().with_u32_structural_encodings(),
770            blob_metadata,
771        )
772        .await;
773    }
774
775    #[tokio::test]
776    async fn test_blob_v2_packed_round_trip() {
777        let blob_metadata = HashMap::from([(
778            lance_arrow::ARROW_EXT_NAME_KEY.to_string(),
779            lance_arrow::BLOB_V2_EXT_NAME.to_string(),
780        )]);
781
782        let kind_field = Arc::new(ArrowField::new("kind", DataType::UInt8, true));
783        let data_field = Arc::new(ArrowField::new("data", DataType::LargeBinary, true));
784        let uri_field = Arc::new(ArrowField::new("uri", DataType::Utf8, true));
785        let blob_id_field = Arc::new(ArrowField::new("blob_id", DataType::UInt32, true));
786        let blob_size_field = Arc::new(ArrowField::new("blob_size", DataType::UInt64, true));
787        let position_field = Arc::new(ArrowField::new("position", DataType::UInt64, true));
788
789        let kind_array = UInt8Array::from(vec![BlobKind::Packed as u8]);
790        let data_array = LargeBinaryArray::from(vec![None::<&[u8]>]);
791        let uri_array = StringArray::from(vec![None::<&str>]);
792        let blob_id_array = UInt32Array::from(vec![7]);
793        let blob_size_array = UInt64Array::from(vec![5]);
794        let position_array = UInt64Array::from(vec![10]);
795
796        let struct_array = StructArray::from(vec![
797            (kind_field, Arc::new(kind_array) as ArrayRef),
798            (data_field, Arc::new(data_array) as ArrayRef),
799            (uri_field, Arc::new(uri_array) as ArrayRef),
800            (blob_id_field, Arc::new(blob_id_array) as ArrayRef),
801            (blob_size_field, Arc::new(blob_size_array) as ArrayRef),
802            (position_field, Arc::new(position_array) as ArrayRef),
803        ]);
804
805        let expected_descriptor = StructArray::from(vec![
806            (
807                Arc::new(ArrowField::new("kind", DataType::UInt8, false)),
808                Arc::new(UInt8Array::from(vec![BlobKind::Packed as u8])) as ArrayRef,
809            ),
810            (
811                Arc::new(ArrowField::new("position", DataType::UInt64, false)),
812                Arc::new(UInt64Array::from(vec![10])) as ArrayRef,
813            ),
814            (
815                Arc::new(ArrowField::new("size", DataType::UInt64, false)),
816                Arc::new(UInt64Array::from(vec![5])) as ArrayRef,
817            ),
818            (
819                Arc::new(ArrowField::new("blob_id", DataType::UInt32, false)),
820                Arc::new(UInt32Array::from(vec![7])) as ArrayRef,
821            ),
822            (
823                Arc::new(ArrowField::new("blob_uri", DataType::Utf8, false)),
824                Arc::new(StringArray::from(vec![""])) as ArrayRef,
825            ),
826        ]);
827
828        check_round_trip_encoding_of_data_with_expected(
829            vec![Arc::new(struct_array)],
830            Some(Arc::new(expected_descriptor)),
831            &TestCases::default().with_u32_structural_encodings(),
832            blob_metadata,
833        )
834        .await;
835    }
836}