1use std::sync::Arc;
19
20use futures::stream::BoxStream;
21use serde::{Deserialize, Serialize};
22use typed_builder::TypedBuilder;
23
24use crate::error::invalid_data;
25use crate::expr::BoundPredicate;
26use crate::spec::{
27 DataContentType, DataFileFormat, ManifestEntryRef, NameMapping, PartitionSpec, Schema,
28 SchemaRef, SortOrderRef, Struct, StructType,
29};
30use crate::{Error, Result};
31
32pub type FileScanTaskStream = BoxStream<'static, Result<FileScanTask>>;
34
35#[derive(Debug, Clone, Deserialize, PartialEq, TypedBuilder)]
37#[serde(try_from = "crate::scan::task::_serde::FileScanTaskSerde")]
38#[builder(
39 field_defaults(setter(prefix = "with_")),
40 build_method(into = Result<FileScanTask>)
41)]
42pub struct FileScanTask {
43 file_size_in_bytes: u64,
46 start: u64,
48 length: u64,
50 #[builder(default)]
55 record_count: Option<u64>,
56
57 #[builder(default)]
62 first_row_id: Option<i64>,
63
64 #[builder(default)]
71 data_sequence_number: Option<i64>,
72
73 data_file_path: String,
75
76 data_file_format: DataFileFormat,
78
79 schema: SchemaRef,
81 project_field_ids: Vec<i32>,
83 #[builder(default)]
85 predicate: Option<BoundPredicate>,
86
87 #[builder(default)]
89 deletes: Vec<FileScanTaskDeleteFile>,
90
91 #[builder(default)]
95 partition: Option<Struct>,
96
97 #[builder(default)]
101 partition_spec: Option<Arc<PartitionSpec>>,
102
103 #[builder(default)]
107 name_mapping: Option<Arc<NameMapping>>,
108
109 #[builder(default)]
118 unified_partition_type: Option<Arc<StructType>>,
119
120 #[builder(default)]
129 sort_order_id: Option<i32>,
130
131 #[builder(default)]
138 sort_order: Option<SortOrderRef>,
139
140 case_sensitive: bool,
142
143 #[builder(default)]
152 key_metadata: Option<Box<[u8]>>,
153}
154
155impl FileScanTask {
156 pub fn file_size_in_bytes(&self) -> u64 {
158 self.file_size_in_bytes
159 }
160
161 pub fn start(&self) -> u64 {
163 self.start
164 }
165
166 pub fn length(&self) -> u64 {
168 self.length
169 }
170
171 pub fn record_count(&self) -> Option<u64> {
173 self.record_count
174 }
175
176 pub fn first_row_id(&self) -> Option<i64> {
178 self.first_row_id
179 }
180
181 pub fn data_sequence_number(&self) -> Option<i64> {
183 self.data_sequence_number
184 }
185
186 pub fn data_file_path(&self) -> &str {
188 &self.data_file_path
189 }
190
191 pub fn data_file_format(&self) -> DataFileFormat {
193 self.data_file_format
194 }
195
196 pub fn schema(&self) -> &Schema {
198 &self.schema
199 }
200
201 pub fn schema_ref(&self) -> SchemaRef {
203 self.schema.clone()
204 }
205
206 pub fn project_field_ids(&self) -> &[i32] {
208 &self.project_field_ids
209 }
210
211 pub fn predicate(&self) -> Option<&BoundPredicate> {
213 self.predicate.as_ref()
214 }
215
216 pub(crate) fn clear_predicate(&mut self) {
223 self.predicate = None;
224 }
225
226 pub fn deletes(&self) -> &[FileScanTaskDeleteFile] {
228 &self.deletes
229 }
230
231 pub fn partition(&self) -> Option<&Struct> {
233 self.partition.as_ref()
234 }
235
236 pub fn partition_spec(&self) -> Option<&Arc<PartitionSpec>> {
238 self.partition_spec.as_ref()
239 }
240
241 pub fn name_mapping(&self) -> Option<&Arc<NameMapping>> {
243 self.name_mapping.as_ref()
244 }
245
246 pub fn unified_partition_type(&self) -> Option<&Arc<StructType>> {
248 self.unified_partition_type.as_ref()
249 }
250
251 pub fn sort_order_id(&self) -> Option<i32> {
253 self.sort_order_id
254 }
255
256 pub fn sort_order(&self) -> Option<&SortOrderRef> {
258 self.sort_order.as_ref()
259 }
260
261 pub fn case_sensitive(&self) -> bool {
263 self.case_sensitive
264 }
265
266 pub fn key_metadata(&self) -> Option<&[u8]> {
268 self.key_metadata.as_deref()
269 }
270
271 fn validate(&self) -> Result<()> {
272 match (self.partition.as_ref(), self.partition_spec.as_deref()) {
273 (None, None) => Ok(()),
274 (None, Some(partition_spec)) if partition_spec.is_unpartitioned() => Ok(()),
275 (None, Some(_)) => Err(invalid_data!(
276 "FileScanTask with a partitioned spec requires partition values"
277 )),
278 (Some(partition), None) if partition.fields().is_empty() => Ok(()),
279 (Some(_), None) => Err(invalid_data!(
280 "Non-empty FileScanTask partition requires a partition spec"
281 )),
282 (Some(partition), Some(partition_spec))
283 if partition.fields().len() != partition_spec.fields().len() =>
284 {
285 Err(invalid_data!(
286 "FileScanTask partition has {} fields but partition spec has {} fields",
287 partition.fields().len(),
288 partition_spec.fields().len()
289 ))
290 }
291 (Some(_), Some(partition_spec)) => {
292 partition_spec.partition_type(&self.schema)?;
293 Ok(())
294 }
295 }
296 }
297}
298
299impl From<FileScanTask> for Result<FileScanTask> {
300 fn from(task: FileScanTask) -> Self {
301 task.validate()?;
302 Ok(task)
303 }
304}
305
306#[derive(Debug)]
307pub(crate) struct DeleteFileContext {
308 pub(crate) manifest_entry: ManifestEntryRef,
309 pub(crate) partition_spec_id: i32,
310}
311
312impl TryFrom<&DeleteFileContext> for FileScanTaskDeleteFile {
313 type Error = Error;
314
315 fn try_from(ctx: &DeleteFileContext) -> Result<Self> {
316 let file_path = ctx.manifest_entry.file_path();
320 let to_offset = |value: Option<i64>, field: &str| -> Result<Option<u64>> {
321 value
322 .map(|value| {
323 u64::try_from(value).map_err(|_| {
324 invalid_data!("delete file {file_path} has negative {field} {value}")
325 })
326 })
327 .transpose()
328 };
329
330 FileScanTaskDeleteFile::builder()
331 .with_file_path(ctx.manifest_entry.file_path().to_string())
332 .with_file_size_in_bytes(ctx.manifest_entry.file_size_in_bytes())
333 .with_file_type(ctx.manifest_entry.content_type())
334 .with_file_format(ctx.manifest_entry.data_file().file_format())
335 .with_partition_spec_id(ctx.partition_spec_id)
336 .with_equality_ids(ctx.manifest_entry.data_file.equality_ids.clone())
337 .with_referenced_data_file(ctx.manifest_entry.data_file.referenced_data_file.clone())
338 .with_content_offset(to_offset(
339 ctx.manifest_entry.data_file.content_offset,
340 "content_offset",
341 )?)
342 .with_content_size_in_bytes(to_offset(
343 ctx.manifest_entry.data_file.content_size_in_bytes,
344 "content_size_in_bytes",
345 )?)
346 .with_record_count(Some(ctx.manifest_entry.record_count()))
347 .with_key_metadata(
348 ctx.manifest_entry
349 .data_file
350 .key_metadata
351 .as_deref()
352 .map(Box::from),
353 )
354 .build()
355 }
356}
357
358#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, TypedBuilder)]
360#[serde(try_from = "crate::scan::task::_serde::FileScanTaskDeleteFileSerde")]
361#[builder(
362 field_defaults(setter(prefix = "with_")),
363 build_method(into = Result<FileScanTaskDeleteFile>)
364)]
365pub struct FileScanTaskDeleteFile {
366 file_path: String,
368
369 file_size_in_bytes: u64,
371
372 file_type: DataContentType,
374
375 file_format: DataFileFormat,
378
379 partition_spec_id: i32,
381
382 #[builder(default)]
384 equality_ids: Option<Vec<i32>>,
385
386 #[serde(default)]
389 #[serde(skip_serializing_if = "Option::is_none")]
390 #[builder(default)]
391 referenced_data_file: Option<String>,
392
393 #[serde(default)]
398 #[serde(skip_serializing_if = "Option::is_none")]
399 #[builder(default)]
400 content_offset: Option<u64>,
401
402 #[serde(default)]
406 #[serde(skip_serializing_if = "Option::is_none")]
407 #[builder(default)]
408 content_size_in_bytes: Option<u64>,
409
410 #[serde(default)]
413 #[serde(skip_serializing_if = "Option::is_none")]
414 #[builder(default)]
415 record_count: Option<u64>,
416
417 #[serde(default)]
427 #[serde(skip_serializing_if = "Option::is_none")]
428 #[builder(default)]
429 key_metadata: Option<Box<[u8]>>,
430}
431
432mod _serde {
433 use std::sync::Arc;
434
435 use serde::{Deserialize, Serialize};
436
437 use super::{FileScanTask, FileScanTaskDeleteFile};
438 use crate::error::invalid_data;
439 use crate::expr::BoundPredicate;
440 use crate::spec::{
441 DataContentType, DataFileFormat, Literal, NameMapping, PartitionSpec, RawLiteral,
442 SchemaRef, SortOrderRef, StructType, Type,
443 };
444 use crate::{Error, Result};
445
446 #[derive(Deserialize)]
447 pub(super) struct FileScanTaskSerde {
448 file_size_in_bytes: u64,
449 start: u64,
450 length: u64,
451 record_count: Option<u64>,
452 first_row_id: Option<i64>,
453 data_sequence_number: Option<i64>,
454 data_file_path: String,
455 data_file_format: DataFileFormat,
456 schema: SchemaRef,
457 project_field_ids: Vec<i32>,
458 predicate: Option<BoundPredicate>,
459 deletes: Vec<FileScanTaskDeleteFile>,
460 #[serde(default)]
461 partition: Option<RawLiteral>,
462 #[serde(default)]
463 partition_spec: Option<Arc<PartitionSpec>>,
464 #[serde(default)]
465 name_mapping: Option<Arc<NameMapping>>,
466 #[serde(default)]
467 unified_partition_type: Option<Arc<StructType>>,
468 #[serde(default)]
469 sort_order_id: Option<i32>,
470 #[serde(default)]
471 sort_order: Option<SortOrderRef>,
472 case_sensitive: bool,
473 #[serde(default)]
474 key_metadata: Option<Box<[u8]>>,
475 }
476
477 #[derive(Serialize)]
478 struct FileScanTaskRefSerde<'a> {
479 file_size_in_bytes: u64,
480 start: u64,
481 length: u64,
482 #[serde(skip_serializing_if = "Option::is_none")]
483 record_count: Option<u64>,
484 #[serde(skip_serializing_if = "Option::is_none")]
485 first_row_id: Option<i64>,
486 #[serde(skip_serializing_if = "Option::is_none")]
487 data_sequence_number: Option<i64>,
488 data_file_path: &'a str,
489 data_file_format: DataFileFormat,
490 schema: &'a SchemaRef,
491 project_field_ids: &'a [i32],
492 #[serde(skip_serializing_if = "Option::is_none")]
493 predicate: Option<&'a BoundPredicate>,
494 deletes: &'a [FileScanTaskDeleteFile],
495 #[serde(skip_serializing_if = "Option::is_none")]
496 partition: Option<RawLiteral>,
497 #[serde(skip_serializing_if = "Option::is_none")]
498 partition_spec: Option<&'a Arc<PartitionSpec>>,
499 #[serde(skip_serializing_if = "Option::is_none")]
500 name_mapping: Option<&'a Arc<NameMapping>>,
501 #[serde(skip_serializing_if = "Option::is_none")]
502 unified_partition_type: Option<&'a Arc<StructType>>,
503 #[serde(skip_serializing_if = "Option::is_none")]
504 sort_order_id: Option<i32>,
505 #[serde(skip_serializing_if = "Option::is_none")]
506 sort_order: Option<&'a SortOrderRef>,
507 case_sensitive: bool,
508 #[serde(skip_serializing_if = "Option::is_none")]
509 key_metadata: Option<&'a [u8]>,
510 }
511
512 fn partition_type(
513 partition_spec: Option<&PartitionSpec>,
514 schema: &crate::spec::Schema,
515 ) -> Result<Type> {
516 let partition_type = match partition_spec {
517 Some(partition_spec) => partition_spec.partition_type(schema)?,
518 None => PartitionSpec::unpartition_spec().partition_type(schema)?,
519 };
520 Ok(Type::Struct(partition_type))
521 }
522
523 impl<'a> TryFrom<&'a FileScanTask> for FileScanTaskRefSerde<'a> {
524 type Error = Error;
525
526 fn try_from(value: &'a FileScanTask) -> Result<Self> {
527 let partition = value
528 .partition
529 .as_ref()
530 .map(|partition| {
531 let partition_type =
532 partition_type(value.partition_spec.as_deref(), &value.schema)?;
533 RawLiteral::try_from(Literal::Struct(partition.clone()), &partition_type)
534 })
535 .transpose()?;
536
537 Ok(Self {
538 file_size_in_bytes: value.file_size_in_bytes,
539 start: value.start,
540 length: value.length,
541 record_count: value.record_count,
542 first_row_id: value.first_row_id,
543 data_sequence_number: value.data_sequence_number,
544 data_file_path: &value.data_file_path,
545 data_file_format: value.data_file_format,
546 schema: &value.schema,
547 project_field_ids: &value.project_field_ids,
548 predicate: value.predicate.as_ref(),
549 deletes: &value.deletes,
550 partition,
551 partition_spec: value.partition_spec.as_ref(),
552 name_mapping: value.name_mapping.as_ref(),
553 unified_partition_type: value.unified_partition_type.as_ref(),
554 sort_order_id: value.sort_order_id,
555 sort_order: value.sort_order.as_ref(),
556 case_sensitive: value.case_sensitive,
557 key_metadata: value.key_metadata.as_deref(),
558 })
559 }
560 }
561
562 impl Serialize for FileScanTask {
563 fn serialize<S>(&self, serializer: S) -> std::result::Result<S::Ok, S::Error>
564 where S: serde::Serializer {
565 FileScanTaskRefSerde::try_from(self)
566 .map_err(serde::ser::Error::custom)?
567 .serialize(serializer)
568 }
569 }
570
571 impl TryFrom<FileScanTaskSerde> for FileScanTask {
572 type Error = Error;
573
574 fn try_from(value: FileScanTaskSerde) -> Result<Self> {
575 let partition = value
576 .partition
577 .map(|partition| {
578 let partition_type =
579 partition_type(value.partition_spec.as_deref(), &value.schema)?;
580 match partition.try_into(&partition_type)? {
581 Some(Literal::Struct(partition)) => Ok(partition),
582 _ => Err(invalid_data!("FileScanTask partition must be a struct")),
583 }
584 })
585 .transpose()?;
586
587 Self::builder()
588 .with_file_size_in_bytes(value.file_size_in_bytes)
589 .with_start(value.start)
590 .with_length(value.length)
591 .with_record_count(value.record_count)
592 .with_first_row_id(value.first_row_id)
593 .with_data_sequence_number(value.data_sequence_number)
594 .with_data_file_path(value.data_file_path)
595 .with_data_file_format(value.data_file_format)
596 .with_schema(value.schema)
597 .with_project_field_ids(value.project_field_ids)
598 .with_predicate(value.predicate)
599 .with_deletes(value.deletes)
600 .with_partition(partition)
601 .with_partition_spec(value.partition_spec)
602 .with_name_mapping(value.name_mapping)
603 .with_unified_partition_type(value.unified_partition_type)
604 .with_sort_order_id(value.sort_order_id)
605 .with_sort_order(value.sort_order)
606 .with_case_sensitive(value.case_sensitive)
607 .with_key_metadata(value.key_metadata)
608 .build()
609 }
610 }
611
612 #[derive(Deserialize)]
613 pub(super) struct FileScanTaskDeleteFileSerde {
614 file_path: String,
615 file_size_in_bytes: u64,
616 file_type: DataContentType,
617 file_format: DataFileFormat,
618 partition_spec_id: i32,
619 equality_ids: Option<Vec<i32>>,
620 #[serde(default)]
621 referenced_data_file: Option<String>,
622 #[serde(default)]
623 content_offset: Option<u64>,
624 #[serde(default)]
625 content_size_in_bytes: Option<u64>,
626 #[serde(default)]
627 record_count: Option<u64>,
628 #[serde(default)]
629 key_metadata: Option<Box<[u8]>>,
630 }
631
632 impl TryFrom<FileScanTaskDeleteFileSerde> for FileScanTaskDeleteFile {
633 type Error = Error;
634
635 fn try_from(value: FileScanTaskDeleteFileSerde) -> Result<Self> {
636 Self::builder()
637 .with_file_path(value.file_path)
638 .with_file_size_in_bytes(value.file_size_in_bytes)
639 .with_file_type(value.file_type)
640 .with_file_format(value.file_format)
641 .with_partition_spec_id(value.partition_spec_id)
642 .with_equality_ids(value.equality_ids)
643 .with_referenced_data_file(value.referenced_data_file)
644 .with_content_offset(value.content_offset)
645 .with_content_size_in_bytes(value.content_size_in_bytes)
646 .with_record_count(value.record_count)
647 .with_key_metadata(value.key_metadata)
648 .build()
649 }
650 }
651}
652
653impl FileScanTaskDeleteFile {
654 pub fn file_path(&self) -> &str {
656 &self.file_path
657 }
658
659 pub fn file_size_in_bytes(&self) -> u64 {
661 self.file_size_in_bytes
662 }
663
664 pub fn file_type(&self) -> DataContentType {
666 self.file_type
667 }
668
669 pub fn file_format(&self) -> DataFileFormat {
671 self.file_format
672 }
673
674 pub fn partition_spec_id(&self) -> i32 {
676 self.partition_spec_id
677 }
678
679 pub fn equality_ids(&self) -> Option<&[i32]> {
681 self.equality_ids.as_deref()
682 }
683
684 pub fn referenced_data_file(&self) -> Option<&str> {
686 self.referenced_data_file.as_deref()
687 }
688
689 pub fn content_offset(&self) -> Option<u64> {
691 self.content_offset
692 }
693
694 pub fn content_size_in_bytes(&self) -> Option<u64> {
696 self.content_size_in_bytes
697 }
698
699 pub fn record_count(&self) -> Option<u64> {
701 self.record_count
702 }
703
704 pub fn key_metadata(&self) -> Option<&[u8]> {
706 self.key_metadata.as_deref()
707 }
708
709 fn is_deletion_vector(&self) -> bool {
710 self.file_type == DataContentType::PositionDeletes
711 && self.file_format == DataFileFormat::Puffin
712 }
713
714 fn validate(&self) -> Result<()> {
715 if self.content_offset.is_some() && self.content_size_in_bytes.is_none() {
721 return Err(invalid_data!(
722 "delete file {} is missing content_size_in_bytes while content_offset is present",
723 self.file_path
724 ));
725 }
726
727 if !self.is_deletion_vector() {
730 return Ok(());
731 }
732
733 let missing = if self.referenced_data_file.is_none() {
734 Some("referenced_data_file")
735 } else if self.content_offset.is_none() {
736 Some("content_offset")
737 } else if self.content_size_in_bytes.is_none() {
738 Some("content_size_in_bytes")
739 } else if self.record_count.is_none() {
740 Some("record_count")
741 } else {
742 None
743 };
744
745 if let Some(field) = missing {
746 return Err(invalid_data!(
747 "deletion vector {} is missing {field}",
748 self.file_path
749 ));
750 }
751
752 Ok(())
753 }
754}
755
756impl From<FileScanTaskDeleteFile> for Result<FileScanTaskDeleteFile> {
757 fn from(task: FileScanTaskDeleteFile) -> Self {
758 task.validate()?;
759 Ok(task)
760 }
761}
762
763#[cfg(test)]
764mod tests {
765 use std::collections::HashMap;
766
767 use super::*;
768 use crate::ErrorKind;
769 use crate::spec::{
770 DataFileBuilder, Literal, ManifestEntry, ManifestStatus, NestedField, PrimitiveType,
771 Transform, Type,
772 };
773
774 fn build_file_scan_task(
775 schema: SchemaRef,
776 partition: Option<Struct>,
777 partition_spec: Option<Arc<PartitionSpec>>,
778 ) -> Result<FileScanTask> {
779 FileScanTask::builder()
780 .with_file_size_in_bytes(100)
781 .with_start(0)
782 .with_length(100)
783 .with_data_file_path("data_file_path".to_string())
784 .with_data_file_format(DataFileFormat::Parquet)
785 .with_schema(schema)
786 .with_project_field_ids(vec![])
787 .with_partition(partition)
788 .with_partition_spec(partition_spec)
789 .with_case_sensitive(false)
790 .build()
791 }
792
793 fn schema_and_spec(
794 primitive_type: PrimitiveType,
795 transform: Transform,
796 ) -> (SchemaRef, Arc<PartitionSpec>) {
797 let schema = Arc::new(
798 Schema::builder()
799 .with_fields(vec![Arc::new(NestedField::required(
800 1,
801 "x",
802 Type::Primitive(primitive_type),
803 ))])
804 .build()
805 .unwrap(),
806 );
807 let partition_spec = Arc::new(
808 PartitionSpec::builder(schema.clone())
809 .add_partition_field("x", "x_partition", transform)
810 .unwrap()
811 .build()
812 .unwrap(),
813 );
814 (schema, partition_spec)
815 }
816
817 fn build_delete_file_task(
818 file_type: DataContentType,
819 file_format: DataFileFormat,
820 ) -> Result<FileScanTaskDeleteFile> {
821 FileScanTaskDeleteFile::builder()
822 .with_file_path("delete-file".to_string())
823 .with_file_size_in_bytes(100)
824 .with_file_type(file_type)
825 .with_file_format(file_format)
826 .with_partition_spec_id(0)
827 .build()
828 }
829
830 fn assert_delete_file_builder_error(
831 result: Result<FileScanTaskDeleteFile>,
832 expected_message: &str,
833 ) {
834 match result {
835 Ok(task) => panic!(
836 "expected delete file builder to fail with `{expected_message}`, but got Ok({task:?})"
837 ),
838 Err(err) => {
839 assert_eq!(err.kind(), ErrorKind::DataInvalid);
840 assert_eq!(err.message(), expected_message);
841 }
842 }
843 }
844
845 #[test]
846 fn test_file_scan_task_builder_rejects_non_empty_partition_without_spec() {
847 let err = build_file_scan_task(
849 Arc::new(Schema::builder().build().unwrap()),
850 Some(Struct::from_iter([Some(Literal::long(42))])),
851 None,
852 )
853 .unwrap_err();
854
855 assert_eq!(err.kind(), ErrorKind::DataInvalid);
856 assert_eq!(
857 err.message(),
858 "Non-empty FileScanTask partition requires a partition spec"
859 );
860 }
861
862 #[test]
863 fn test_file_scan_task_builder_accepts_empty_partition_without_spec() {
864 build_file_scan_task(
865 Arc::new(Schema::builder().build().unwrap()),
866 Some(Struct::empty()),
867 None,
868 )
869 .unwrap();
870 }
871
872 #[test]
873 fn test_file_scan_task_builder_rejects_partitioned_spec_without_partition() {
874 let (schema, partition_spec) = schema_and_spec(PrimitiveType::Long, Transform::Identity);
875
876 let err = build_file_scan_task(schema, None, Some(partition_spec)).unwrap_err();
877
878 assert_eq!(err.kind(), ErrorKind::DataInvalid);
879 assert_eq!(
880 err.message(),
881 "FileScanTask with a partitioned spec requires partition values"
882 );
883 }
884
885 #[test]
886 fn test_file_scan_task_builder_accepts_unpartitioned_spec_without_partition() {
887 build_file_scan_task(
888 Arc::new(Schema::builder().build().unwrap()),
889 None,
890 Some(Arc::new(PartitionSpec::unpartition_spec())),
891 )
892 .unwrap();
893 }
894
895 #[test]
896 fn test_file_scan_task_builder_rejects_partition_arity_mismatch() {
897 let (schema, partition_spec) = schema_and_spec(PrimitiveType::Long, Transform::Identity);
898
899 let err =
900 build_file_scan_task(schema, Some(Struct::empty()), Some(partition_spec)).unwrap_err();
901
902 assert_eq!(err.kind(), ErrorKind::DataInvalid);
903 assert!(err.message().contains("partition has 0 fields"));
904 assert!(err.message().contains("partition spec has 1 fields"));
905 }
906
907 #[test]
908 fn test_file_scan_task_builder_accepts_dropped_partition_source_column() {
909 let (_historical_schema, partition_spec) =
910 schema_and_spec(PrimitiveType::Long, Transform::Identity);
911 let current_schema = Arc::new(
912 Schema::builder()
913 .with_fields(vec![Arc::new(NestedField::required(
914 2,
915 "y",
916 Type::Primitive(PrimitiveType::String),
917 ))])
918 .build()
919 .unwrap(),
920 );
921
922 build_file_scan_task(
923 current_schema,
924 Some(Struct::from_iter([Some(Literal::long(42))])),
925 Some(partition_spec),
926 )
927 .unwrap();
928 }
929
930 #[test]
931 fn test_file_scan_task_builder_rejects_partition_spec_incompatible_with_schema() {
932 let (_historical_schema, partition_spec) =
933 schema_and_spec(PrimitiveType::Timestamp, Transform::Day);
934 let current_schema = Arc::new(
935 Schema::builder()
936 .with_fields(vec![Arc::new(NestedField::required(
937 1,
938 "x",
939 Type::Primitive(PrimitiveType::String),
940 ))])
941 .build()
942 .unwrap(),
943 );
944
945 let err = build_file_scan_task(
946 current_schema,
947 Some(Struct::from_iter([Some(Literal::date(20_000))])),
948 Some(partition_spec),
949 )
950 .unwrap_err();
951
952 assert_eq!(err.kind(), ErrorKind::DataInvalid);
953 }
954
955 #[test]
956 fn test_delete_file_builder_accepts_valid_deletion_vector() {
957 let task = FileScanTaskDeleteFile::builder()
958 .with_file_path("dv.puffin".to_string())
959 .with_file_size_in_bytes(100)
960 .with_file_type(DataContentType::PositionDeletes)
961 .with_file_format(DataFileFormat::Puffin)
962 .with_partition_spec_id(7)
963 .with_referenced_data_file(Some("data.parquet".to_string()))
964 .with_content_offset(Some(11))
965 .with_content_size_in_bytes(Some(13))
966 .with_record_count(Some(3))
967 .with_key_metadata(Some(vec![17, 19].into_boxed_slice()))
968 .build()
969 .unwrap();
970
971 assert_eq!(task.file_path(), "dv.puffin");
972 assert_eq!(task.file_size_in_bytes(), 100);
973 assert_eq!(task.file_type(), DataContentType::PositionDeletes);
974 assert_eq!(task.file_format(), DataFileFormat::Puffin);
975 assert_eq!(task.partition_spec_id(), 7);
976 assert_eq!(task.equality_ids(), None);
977 assert_eq!(task.referenced_data_file(), Some("data.parquet"));
978 assert_eq!(task.content_offset(), Some(11));
979 assert_eq!(task.content_size_in_bytes(), Some(13));
980 assert_eq!(task.record_count(), Some(3));
981 assert_eq!(task.key_metadata(), Some([17, 19].as_slice()));
982 }
983
984 #[test]
985 fn test_delete_file_builder_rejects_dv_missing_referenced_data_file() {
986 assert_delete_file_builder_error(
987 FileScanTaskDeleteFile::builder()
988 .with_file_path("dv.puffin".to_string())
989 .with_file_size_in_bytes(100)
990 .with_file_type(DataContentType::PositionDeletes)
991 .with_file_format(DataFileFormat::Puffin)
992 .with_partition_spec_id(0)
993 .with_content_offset(Some(7))
994 .with_content_size_in_bytes(Some(11))
995 .with_record_count(Some(3))
996 .build(),
997 "deletion vector dv.puffin is missing referenced_data_file",
998 );
999 }
1000
1001 #[test]
1002 fn test_delete_file_builder_rejects_dv_missing_content_offset() {
1003 assert_delete_file_builder_error(
1004 FileScanTaskDeleteFile::builder()
1005 .with_file_path("dv.puffin".to_string())
1006 .with_file_size_in_bytes(100)
1007 .with_file_type(DataContentType::PositionDeletes)
1008 .with_file_format(DataFileFormat::Puffin)
1009 .with_partition_spec_id(0)
1010 .with_referenced_data_file(Some("data.parquet".to_string()))
1011 .with_content_size_in_bytes(Some(11))
1012 .with_record_count(Some(3))
1013 .build(),
1014 "deletion vector dv.puffin is missing content_offset",
1015 );
1016 }
1017
1018 #[test]
1019 fn test_delete_file_builder_rejects_dv_missing_content_size() {
1020 assert_delete_file_builder_error(
1021 FileScanTaskDeleteFile::builder()
1022 .with_file_path("dv.puffin".to_string())
1023 .with_file_size_in_bytes(100)
1024 .with_file_type(DataContentType::PositionDeletes)
1025 .with_file_format(DataFileFormat::Puffin)
1026 .with_partition_spec_id(0)
1027 .with_referenced_data_file(Some("data.parquet".to_string()))
1028 .with_content_offset(Some(7))
1029 .with_record_count(Some(3))
1030 .build(),
1031 "delete file dv.puffin is missing content_size_in_bytes while content_offset is present",
1032 );
1033 }
1034
1035 #[test]
1036 fn test_delete_file_builder_rejects_dv_missing_record_count() {
1037 assert_delete_file_builder_error(
1038 FileScanTaskDeleteFile::builder()
1039 .with_file_path("dv.puffin".to_string())
1040 .with_file_size_in_bytes(100)
1041 .with_file_type(DataContentType::PositionDeletes)
1042 .with_file_format(DataFileFormat::Puffin)
1043 .with_partition_spec_id(0)
1044 .with_referenced_data_file(Some("data.parquet".to_string()))
1045 .with_content_offset(Some(7))
1046 .with_content_size_in_bytes(Some(11))
1047 .build(),
1048 "deletion vector dv.puffin is missing record_count",
1049 );
1050 }
1051
1052 #[test]
1053 fn test_delete_file_builder_rejects_non_dv_offset_without_size() {
1054 assert_delete_file_builder_error(
1058 FileScanTaskDeleteFile::builder()
1059 .with_file_path("position-deletes.parquet".to_string())
1060 .with_file_size_in_bytes(100)
1061 .with_file_type(DataContentType::PositionDeletes)
1062 .with_file_format(DataFileFormat::Parquet)
1063 .with_partition_spec_id(0)
1064 .with_content_offset(Some(7))
1065 .build(),
1066 "delete file position-deletes.parquet is missing content_size_in_bytes while content_offset is present",
1067 );
1068 }
1069
1070 fn delete_manifest_entry(
1071 file_path: &str,
1072 file_format: DataFileFormat,
1073 content_offset: Option<i64>,
1074 content_size_in_bytes: Option<i64>,
1075 ) -> DeleteFileContext {
1076 let mut data_file = DataFileBuilder::default()
1077 .content(DataContentType::PositionDeletes)
1078 .file_path(file_path.to_string())
1079 .file_format(file_format)
1080 .partition(Struct::empty())
1081 .record_count(3)
1082 .file_size_in_bytes(100)
1083 .column_sizes(HashMap::new())
1084 .value_counts(HashMap::new())
1085 .null_value_counts(HashMap::new())
1086 .partition_spec_id(0)
1087 .referenced_data_file(Some("data.parquet".to_string()))
1088 .build()
1089 .unwrap();
1090 data_file.content_offset = content_offset;
1091 data_file.content_size_in_bytes = content_size_in_bytes;
1092
1093 DeleteFileContext {
1094 manifest_entry: Arc::new(
1095 ManifestEntry::builder()
1096 .status(ManifestStatus::Added)
1097 .data_file(data_file)
1098 .build(),
1099 ),
1100 partition_spec_id: 0,
1101 }
1102 }
1103
1104 fn assert_delete_context_error(ctx: &DeleteFileContext, expected_message: &str) {
1105 match FileScanTaskDeleteFile::try_from(ctx) {
1106 Ok(task) => panic!(
1107 "expected conversion to fail with `{expected_message}`, but got Ok({task:?})"
1108 ),
1109 Err(error) => {
1110 assert_eq!(error.kind(), ErrorKind::DataInvalid);
1111 assert!(
1112 error.message().contains(expected_message),
1113 "expected `{expected_message}`, got `{}`",
1114 error.message()
1115 );
1116 }
1117 }
1118 }
1119
1120 #[test]
1124 fn test_delete_context_rejects_negative_dv_offset() {
1125 assert_delete_context_error(
1126 &delete_manifest_entry("dv.puffin", DataFileFormat::Puffin, Some(-1), Some(11)),
1127 "delete file dv.puffin has negative content_offset -1",
1128 );
1129 }
1130
1131 #[test]
1132 fn test_delete_context_rejects_negative_dv_size() {
1133 assert_delete_context_error(
1134 &delete_manifest_entry("dv.puffin", DataFileFormat::Puffin, Some(7), Some(-1)),
1135 "delete file dv.puffin has negative content_size_in_bytes -1",
1136 );
1137 }
1138
1139 #[test]
1142 fn test_delete_context_rejects_negative_coordinates_for_non_dv() {
1143 assert_delete_context_error(
1144 &delete_manifest_entry(
1145 "position-deletes.parquet",
1146 DataFileFormat::Parquet,
1147 Some(-1),
1148 None,
1149 ),
1150 "delete file position-deletes.parquet has negative content_offset -1",
1151 );
1152 assert_delete_context_error(
1153 &delete_manifest_entry(
1154 "position-deletes.parquet",
1155 DataFileFormat::Parquet,
1156 None,
1157 Some(-1),
1158 ),
1159 "delete file position-deletes.parquet has negative content_size_in_bytes -1",
1160 );
1161 }
1162
1163 #[test]
1164 fn test_delete_file_builder_accepts_non_dv_delete_without_dv_fields() {
1165 build_delete_file_task(DataContentType::PositionDeletes, DataFileFormat::Parquet).unwrap();
1166 build_delete_file_task(DataContentType::EqualityDeletes, DataFileFormat::Parquet).unwrap();
1167 }
1168}