1use std::ops::Range;
13use std::sync::Arc;
14
15use arrow_array::LargeBinaryArray;
16use arrow_array::builder::LargeBinaryBuilder;
17use arrow_schema::{DataType, Field, Schema};
18use lance::dataset::{BlobRangeRequest as LanceBlobRangeRequest, Dataset, WriteParams};
19use lance_arrow::FieldExt;
20use lance_file::version::LanceFileVersion;
21use lance_io::object_store::ObjectStore;
22use object_store::path::Path;
23
24use crate::error::{Error, Result};
25
26#[derive(Debug)]
29pub struct BlobFile {
30 inner: BlobFileInner,
31}
32
33#[derive(Debug)]
34enum BlobFileInner {
35 Native(lance::dataset::BlobFile),
36 #[cfg(feature = "remote")]
37 Remote(Box<crate::remote::table::blobs::RemoteBlobFile>),
38}
39
40impl From<lance::dataset::BlobFile> for BlobFile {
41 fn from(value: lance::dataset::BlobFile) -> Self {
42 Self {
43 inner: BlobFileInner::Native(value),
44 }
45 }
46}
47
48#[cfg(feature = "remote")]
49impl From<crate::remote::table::blobs::RemoteBlobFile> for BlobFile {
50 fn from(value: crate::remote::table::blobs::RemoteBlobFile) -> Self {
51 Self {
52 inner: BlobFileInner::Remote(Box::new(value)),
53 }
54 }
55}
56
57impl BlobFile {
58 pub fn new_inline(
60 object_store: Arc<ObjectStore>,
61 path: Path,
62 position: u64,
63 size: u64,
64 ) -> Self {
65 lance::dataset::BlobFile::new_inline(object_store, path, position, size).into()
66 }
67
68 pub fn new_dedicated(object_store: Arc<ObjectStore>, path: Path, size: u64) -> Self {
70 lance::dataset::BlobFile::new_dedicated(object_store, path, size).into()
71 }
72
73 pub fn new_packed(
75 object_store: Arc<ObjectStore>,
76 path: Path,
77 position: u64,
78 size: u64,
79 ) -> Self {
80 lance::dataset::BlobFile::new_packed(object_store, path, position, size).into()
81 }
82
83 pub fn new_external(
85 object_store: Arc<ObjectStore>,
86 path: Path,
87 uri: String,
88 position: u64,
89 size: u64,
90 ) -> Self {
91 lance::dataset::BlobFile::new_external(object_store, path, uri, position, size).into()
92 }
93
94 pub async fn close(&self) -> lance_core::Result<()> {
96 match &self.inner {
97 BlobFileInner::Native(file) => file.close().await,
98 #[cfg(feature = "remote")]
99 BlobFileInner::Remote(file) => file.close().await,
100 }
101 }
102
103 pub async fn is_closed(&self) -> bool {
105 match &self.inner {
106 BlobFileInner::Native(file) => file.is_closed().await,
107 #[cfg(feature = "remote")]
108 BlobFileInner::Remote(file) => file.is_closed(),
109 }
110 }
111
112 pub async fn read_range(&self, range: Range<u64>) -> lance_core::Result<bytes::Bytes> {
114 match &self.inner {
115 BlobFileInner::Native(file) => file.read_range(range).await,
116 #[cfg(feature = "remote")]
117 BlobFileInner::Remote(file) => file.read_range(range).await,
118 }
119 }
120
121 pub async fn read_ranges(
123 &self,
124 ranges: &[Range<u64>],
125 ) -> lance_core::Result<Vec<bytes::Bytes>> {
126 match &self.inner {
127 BlobFileInner::Native(file) => file.read_ranges(ranges).await,
128 #[cfg(feature = "remote")]
129 BlobFileInner::Remote(file) => file.read_ranges(ranges).await,
130 }
131 }
132
133 pub async fn read(&self) -> lance_core::Result<bytes::Bytes> {
135 match &self.inner {
136 BlobFileInner::Native(file) => file.read().await,
137 #[cfg(feature = "remote")]
138 BlobFileInner::Remote(file) => file.read().await,
139 }
140 }
141
142 pub async fn read_up_to(&self, len: usize) -> lance_core::Result<bytes::Bytes> {
144 match &self.inner {
145 BlobFileInner::Native(file) => file.read_up_to(len).await,
146 #[cfg(feature = "remote")]
147 BlobFileInner::Remote(file) => file.read_up_to(len).await,
148 }
149 }
150
151 pub async fn seek(&self, new_cursor: u64) -> lance_core::Result<()> {
153 match &self.inner {
154 BlobFileInner::Native(file) => file.seek(new_cursor).await,
155 #[cfg(feature = "remote")]
156 BlobFileInner::Remote(file) => file.seek(new_cursor).await,
157 }
158 }
159
160 pub async fn tell(&self) -> lance_core::Result<u64> {
162 match &self.inner {
163 BlobFileInner::Native(file) => file.tell().await,
164 #[cfg(feature = "remote")]
165 BlobFileInner::Remote(file) => file.tell().await,
166 }
167 }
168
169 pub fn size(&self) -> u64 {
171 match &self.inner {
172 BlobFileInner::Native(file) => file.size(),
173 #[cfg(feature = "remote")]
174 BlobFileInner::Remote(file) => file.size(),
175 }
176 }
177
178 pub fn position(&self) -> Option<u64> {
181 match &self.inner {
182 BlobFileInner::Native(file) => Some(file.position()),
183 #[cfg(feature = "remote")]
184 BlobFileInner::Remote(_) => None,
185 }
186 }
187
188 pub fn data_path(&self) -> Option<&Path> {
191 match &self.inner {
192 BlobFileInner::Native(file) => Some(file.data_path()),
193 #[cfg(feature = "remote")]
194 BlobFileInner::Remote(_) => None,
195 }
196 }
197
198 pub fn kind(&self) -> Option<lance_core::datatypes::BlobKind> {
201 match &self.inner {
202 BlobFileInner::Native(file) => Some(file.kind()),
203 #[cfg(feature = "remote")]
204 BlobFileInner::Remote(_) => None,
205 }
206 }
207
208 pub fn uri(&self) -> Option<&str> {
210 match &self.inner {
211 BlobFileInner::Native(file) => file.uri(),
212 #[cfg(feature = "remote")]
213 BlobFileInner::Remote(_) => None,
214 }
215 }
216}
217
218#[derive(Debug, Clone, Copy, PartialEq, Eq)]
223pub struct BlobRangeRequest {
224 pub row_id: u64,
226 pub offset: u64,
228 pub length: u64,
230}
231
232impl BlobRangeRequest {
233 pub const fn new(row_id: u64, offset: u64, length: u64) -> Self {
235 Self {
236 row_id,
237 offset,
238 length,
239 }
240 }
241}
242
243pub fn blob(name: impl AsRef<str>, nullable: bool) -> Field {
260 lance::blob::blob_field(name.as_ref(), nullable)
261}
262
263pub fn is_blob(field: &Field) -> bool {
270 field.is_blob_v2()
271}
272
273fn field_tree_has_blob_v2(field: &Field) -> bool {
275 if field.is_blob_v2() {
276 return true;
277 }
278 match field.data_type() {
279 DataType::Struct(children) => children.iter().any(|c| field_tree_has_blob_v2(c)),
280 DataType::List(child) | DataType::LargeList(child) | DataType::FixedSizeList(child, _) => {
281 field_tree_has_blob_v2(child)
282 }
283 _ => false,
284 }
285}
286
287fn collect_blob_paths(field: &Field, prefix: &str, paths: &mut Vec<String>) {
289 let path = if prefix.is_empty() {
290 field.name().clone()
291 } else {
292 format!("{prefix}.{}", field.name())
293 };
294 if field.is_blob_v2() {
295 paths.push(path);
296 return;
297 }
298 match field.data_type() {
299 DataType::Struct(children) => {
300 for child in children {
301 collect_blob_paths(child, &path, paths);
302 }
303 }
304 DataType::List(child) | DataType::LargeList(child) | DataType::FixedSizeList(child, _) => {
305 collect_blob_paths(child, &path, paths)
306 }
307 _ => {}
308 }
309}
310
311pub(crate) fn has_blob_columns(schema: &Schema) -> bool {
313 schema.fields().iter().any(|f| field_tree_has_blob_v2(f))
314}
315
316pub(crate) fn blob_column_names(schema: &Schema) -> Vec<String> {
319 let mut paths = Vec::new();
320 for field in schema.fields() {
321 collect_blob_paths(field, "", &mut paths);
322 }
323 paths
324}
325
326pub(crate) fn ensure_blob_storage_version(schema: &Schema, params: &mut WriteParams) {
328 if !has_blob_columns(schema) {
329 return;
330 }
331
332 let resolved = params
333 .data_storage_version
334 .unwrap_or(LanceFileVersion::Stable)
335 .resolve();
336 if resolved < LanceFileVersion::V2_2 {
337 params.data_storage_version = Some(LanceFileVersion::V2_2);
338 }
339}
340
341pub(crate) fn ensure_blob_v2_column(
345 schema: &lance_core::datatypes::Schema,
346 column: &str,
347) -> Result<()> {
348 match schema.field(column) {
349 Some(field) if field.is_blob_v2() => Ok(()),
350 Some(field) if field.is_blob() => Err(Error::InvalidInput {
351 message: format!(
352 "column '{column}' is a legacy blob column; blob APIs require blob v2 columns \
353 (ARROW:extension:name = \"lance.blob.v2\")"
354 ),
355 }),
356 Some(_) => Err(Error::InvalidInput {
357 message: format!("column '{column}' is not a blob column"),
358 }),
359 None => Err(Error::InvalidInput {
360 message: format!("no column named '{column}' in this table"),
361 }),
362 }
363}
364
365fn ensure_all_row_ids_resolved(column: &str, requested: usize, resolved: usize) -> Result<()> {
366 if requested == resolved {
367 return Ok(());
368 }
369 if resolved < requested {
370 Err(Error::InvalidInput {
371 message: format!(
372 "blob read for column '{column}' requested {requested} row ids but only {resolved} \
373 exist in the table; pass row ids collected from this table"
374 ),
375 })
376 } else {
377 Err(Error::Runtime {
378 message: format!(
379 "blob read for column '{column}' returned {resolved} results for {requested} row ids"
380 ),
381 })
382 }
383}
384
385pub(crate) async fn take_blob_ranges_aligned(
387 dataset: &Arc<Dataset>,
388 column: &str,
389 requests: &[BlobRangeRequest],
390) -> Result<LargeBinaryArray> {
391 ensure_blob_v2_column(dataset.schema(), column)?;
392 if requests.is_empty() {
393 return Ok(LargeBinaryBuilder::new().finish());
394 }
395
396 let lance_requests = requests
397 .iter()
398 .map(|request| LanceBlobRangeRequest::new(request.row_id, request.offset, request.length))
399 .collect::<Vec<_>>();
400 let payloads = dataset
401 .read_blob_ranges(column)?
402 .with_row_ids(lance_requests)
403 .preserve_order(true)
404 .execute()
405 .await?;
406 ensure_all_row_ids_resolved(column, requests.len(), payloads.len())?;
407
408 let mut builder = LargeBinaryBuilder::new();
409 for payload in payloads {
410 match payload.data {
411 Some(data) => builder.append_value(data),
412 None => builder.append_null(),
413 }
414 }
415 Ok(builder.finish())
416}
417
418pub(crate) async fn take_blobs_aligned(
420 dataset: &Arc<Dataset>,
421 column: &str,
422 row_ids: &[u64],
423) -> Result<LargeBinaryArray> {
424 ensure_blob_v2_column(dataset.schema(), column)?;
425 if row_ids.is_empty() {
426 return Ok(LargeBinaryBuilder::new().finish());
427 }
428
429 let payloads = dataset
430 .read_blobs(column)?
431 .with_row_ids(row_ids.to_vec())
432 .preserve_order(true)
433 .execute()
434 .await?;
435 ensure_all_row_ids_resolved(column, row_ids.len(), payloads.len())?;
436
437 let mut builder = LargeBinaryBuilder::new();
438 for payload in payloads {
439 match payload.data {
440 Some(data) => builder.append_value(data),
441 None => builder.append_null(),
442 }
443 }
444 Ok(builder.finish())
445}
446
447pub(crate) async fn take_blob_files_aligned(
449 dataset: &Arc<Dataset>,
450 column: &str,
451 row_ids: &[u64],
452) -> Result<Vec<Option<BlobFile>>> {
453 ensure_blob_v2_column(dataset.schema(), column)?;
454 if row_ids.is_empty() {
455 return Ok(Vec::new());
456 }
457
458 let handles = dataset.take_blobs(row_ids, column).await?;
459 ensure_all_row_ids_resolved(column, row_ids.len(), handles.len())?;
460 Ok(handles
461 .into_iter()
462 .map(|handle| handle.map(Into::into))
463 .collect())
464}
465
466#[cfg(test)]
467mod tests {
468 use super::*;
469 use arrow_schema::DataType;
470 use lance_arrow::ARROW_EXT_NAME_KEY;
471
472 fn blob_schema() -> Schema {
473 Schema::new(vec![
474 Field::new("id", DataType::Int64, false),
475 blob("image", true),
476 ])
477 }
478
479 #[test]
480 fn blob_field_carries_v2_extension_marker() {
481 let field = blob("image", true);
482 assert_eq!(
483 field.metadata().get(ARROW_EXT_NAME_KEY).map(String::as_str),
484 Some("lance.blob.v2")
485 );
486 assert!(matches!(field.data_type(), DataType::Struct(_)));
487 }
488
489 #[test]
490 fn has_blob_columns_detects_blob_fields() {
491 assert!(has_blob_columns(&blob_schema()));
492 let plain = Schema::new(vec![Field::new("id", DataType::Int64, false)]);
493 assert!(!has_blob_columns(&plain));
494 }
495
496 #[test]
497 fn storage_version_bumps_to_v2_2() {
498 let mut params = WriteParams::default();
499 ensure_blob_storage_version(&blob_schema(), &mut params);
500 assert_eq!(
501 params.data_storage_version.unwrap().resolve(),
502 LanceFileVersion::V2_2
503 );
504 }
505
506 #[test]
507 fn storage_version_overrides_lower_explicit_version() {
508 let mut params = WriteParams {
509 data_storage_version: Some(LanceFileVersion::V2_0),
510 ..Default::default()
511 };
512 ensure_blob_storage_version(&blob_schema(), &mut params);
513 assert_eq!(
514 params.data_storage_version.unwrap().resolve(),
515 LanceFileVersion::V2_2
516 );
517 }
518
519 #[test]
520 fn storage_version_keeps_higher_explicit_version() {
521 let mut params = WriteParams {
522 data_storage_version: Some(LanceFileVersion::V2_3),
523 ..Default::default()
524 };
525 ensure_blob_storage_version(&blob_schema(), &mut params);
526 assert_eq!(params.data_storage_version.unwrap(), LanceFileVersion::V2_3);
527 }
528
529 #[test]
530 fn legacy_v1_blob_column_is_rejected_with_migration_hint() {
531 let legacy = Field::new("image", DataType::LargeBinary, true).with_metadata(
532 std::collections::HashMap::from([(
533 "lance-encoding:blob".to_string(),
534 "true".to_string(),
535 )]),
536 );
537 let arrow_schema = Schema::new(vec![legacy]);
538 let lance_schema = lance_core::datatypes::Schema::try_from(&arrow_schema).unwrap();
539
540 let err = ensure_blob_v2_column(&lance_schema, "image").unwrap_err();
541 assert!(matches!(err, Error::InvalidInput { .. }));
542 assert!(err.to_string().contains("legacy blob column"));
543 assert!(err.to_string().contains("lance.blob.v2"));
544 }
545
546 #[test]
547 fn non_blob_and_unknown_columns_are_rejected_by_name() {
548 let arrow_schema = Schema::new(vec![Field::new("id", DataType::Int64, false)]);
549 let lance_schema = lance_core::datatypes::Schema::try_from(&arrow_schema).unwrap();
550
551 let err = ensure_blob_v2_column(&lance_schema, "id").unwrap_err();
552 assert!(err.to_string().contains("'id' is not a blob column"));
553
554 let err = ensure_blob_v2_column(&lance_schema, "missing").unwrap_err();
555 assert!(err.to_string().contains("no column named 'missing'"));
556 }
557
558 #[test]
559 fn blob_column_names_includes_nested_path() {
560 let blob_field = blob("blob", true);
561 let info = Field::new(
562 "info",
563 DataType::Struct(vec![Field::new("name", DataType::Utf8, false), blob_field].into()),
564 true,
565 );
566 let schema = Schema::new(vec![Field::new("id", DataType::Int64, false), info]);
567 assert_eq!(blob_column_names(&schema), vec!["info.blob"]);
568 }
569
570 #[test]
571 fn storage_version_noop_without_blob_columns() {
572 let schema = Schema::new(vec![Field::new("id", DataType::Int64, false)]);
573 let mut params = WriteParams::default();
574 ensure_blob_storage_version(&schema, &mut params);
575 assert!(params.data_storage_version.is_none());
576 }
577}