Skip to main content

delta_arrow_reader/
reader.rs

1//! Public Delta-to-Arrow reader and optional DataFusion adapter.
2
3pub(crate) mod backend;
4#[cfg(feature = "datafusion")]
5pub mod datafusion;
6pub(crate) mod deletion_vector;
7pub(crate) mod metrics;
8mod options;
9pub(crate) mod partition_target;
10pub(crate) mod planning;
11pub(crate) mod predicate;
12#[allow(dead_code)]
13pub(crate) mod scheduling;
14pub(crate) mod transform;
15
16#[doc(hidden)]
17pub use metrics::ParquetRangePlanningDiagnosticSnapshot;
18pub use metrics::{DeltaScanMetrics, DeltaScanMetricsSnapshot};
19#[doc(hidden)]
20pub use options::ParquetRangeReadPolicy;
21pub use options::{
22    DeltaScanExecutionOptions, DeltaSnapshotSelection, DeltaStorageOptions, ParquetReaderBackend,
23};
24pub use predicate::{DeltaComparison, DeltaPredicate, DeltaScalar};
25
26use std::{
27    collections::VecDeque,
28    fmt,
29    pin::Pin,
30    sync::Arc,
31    task::{Context, Poll},
32};
33
34use arrow::{datatypes::SchemaRef, record_batch::RecordBatch};
35use futures_util::Stream;
36use snafu::ResultExt;
37
38use self::{
39    planning::{DeltaScanPartitionTargetOptions, DeltaScanPlan, plan_scan},
40    predicate::{evaluate_predicate, referenced_columns, validate_predicate},
41    scheduling::{
42        DeltaScanScheduler, FileAdmissionDecision, FileAdmissionPolicy, FileBatchStream,
43        FileExecutor, PartitionStream,
44    },
45};
46
47use crate::{
48    DeltaProtocol, DeltaReaderError,
49    delta::{
50        kernel::{kernel_pruning_is_exact, kernel_pruning_predicate},
51        protocol::validate_protocol,
52        snapshot::{
53            ArrowTableSnapshot, KernelTableSnapshot, load_delta_table_snapshot,
54            load_kernel_table_snapshot,
55        },
56    },
57    error::{DataFileReadSnafu, InvalidConfigurationSnafu, ScanPlanningSnafu, SnapshotLoadSnafu},
58};
59
60const TRACING_TARGET: &str = "delta_arrow_reader";
61const DELTA_LOG_SCAN_METADATA_SOURCE: &str = "delta_log";
62const QUERY_PLANNING_CACHE_SCAN_METADATA_SOURCE: &str = "query_planning_cache";
63
64/// Selects how much work table loading performs before returning.
65///
66/// Each warmup level includes the work performed by the levels before it. The default performs no
67/// warmup and leaves reusable metadata work to individual queries.
68#[non_exhaustive]
69#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
70pub enum WarmupMode {
71    /// Performs no warmup.
72    #[default]
73    None,
74    /// Loads and retains the active-file metadata used during query planning.
75    QueryPlanning,
76}
77
78/// Configures and loads one immutable Delta table snapshot.
79///
80/// The asynchronous path uses the caller's Tokio runtime. Scans return a
81/// pull-driven stream and do not materialize the whole table.
82///
83/// # Example
84///
85/// ```no_run
86/// use delta_arrow_reader::{DeltaComparison, DeltaPredicate, DeltaScalar, DeltaTableBuilder};
87/// use futures_util::TryStreamExt;
88///
89/// # async fn read_table() -> Result<(), Box<dyn std::error::Error>> {
90/// let table = DeltaTableBuilder::new("/tmp/example-delta-table")
91///     .load_table()
92///     .await?;
93/// let scan = table
94///     .scan()
95///     .with_projection(["id", "name"])
96///     .with_predicate(DeltaPredicate::Compare {
97///         column: "id".into(),
98///         op: DeltaComparison::GtEq,
99///         value: DeltaScalar::Int64(10),
100///     })
101///     .with_limit(100)
102///     .build()
103///     .await?;
104/// let mut batches = scan.into_stream();
105///
106/// while let Some(batch) = batches.try_next().await? {
107///     println!("rows={}", batch.num_rows());
108/// }
109/// # Ok(())
110/// # }
111/// ```
112#[must_use = "table builder settings do nothing unless the builder is loaded"]
113pub struct DeltaTableBuilder {
114    table_location: String,
115    storage_options: DeltaStorageOptions,
116    snapshot_selection: DeltaSnapshotSelection,
117    execution_options: DeltaScanExecutionOptions,
118    warmup: WarmupMode,
119}
120
121impl DeltaTableBuilder {
122    /// Creates a builder for the latest snapshot with default execution settings.
123    pub fn new(table_location: impl Into<String>) -> Self {
124        Self {
125            table_location: table_location.into(),
126            storage_options: DeltaStorageOptions::new(),
127            snapshot_selection: DeltaSnapshotSelection::Latest,
128            execution_options: DeltaScanExecutionOptions::new(),
129            warmup: WarmupMode::None,
130        }
131    }
132
133    /// Replaces the storage options forwarded during table loading.
134    pub fn with_storage_options(mut self, storage_options: DeltaStorageOptions) -> Self {
135        self.storage_options = storage_options;
136        self
137    }
138
139    /// Selects the Delta snapshot to load.
140    pub const fn with_snapshot_selection(
141        mut self,
142        snapshot_selection: DeltaSnapshotSelection,
143    ) -> Self {
144        self.snapshot_selection = snapshot_selection;
145        self
146    }
147
148    /// Replaces the default execution settings used by scans of this table.
149    pub const fn with_execution_options(
150        mut self,
151        execution_options: DeltaScanExecutionOptions,
152    ) -> Self {
153        self.execution_options = execution_options;
154        self
155    }
156
157    /// Selects work to finish during [`Self::load_table`] for reuse by later queries.
158    ///
159    /// [`Self::load_snapshot`] does not build a table and rejects any setting other than
160    /// [`WarmupMode::None`].
161    pub const fn with_warmup(mut self, warmup: WarmupMode) -> Self {
162        self.warmup = warmup;
163        self
164    }
165
166    /// Loads the table through the caller-owned Tokio runtime.
167    ///
168    /// With [`WarmupMode::None`], this method loads the snapshot and converts its schema. Protocol
169    /// validation and scan metadata loading remain deferred until a scan is built. This allows an
170    /// application to inspect the version, schema, and protocol of a table that the reader cannot
171    /// scan.
172    ///
173    /// A warmup mode validates the protocol while loading because it builds reusable query-planning
174    /// state. Query-planning warmup retains active-file metadata.
175    ///
176    /// # Errors
177    ///
178    /// Returns an error if the table location, storage, snapshot, or schema cannot be loaded. A
179    /// warmup can also fail if the table protocol is unsupported or its reusable metadata cannot be
180    /// built.
181    pub async fn load_table(self) -> Result<DeltaTable, DeltaReaderError> {
182        let snapshot = load_delta_table_snapshot(
183            self.table_location,
184            self.storage_options,
185            self.snapshot_selection,
186        )
187        .await?;
188        let snapshot = if self.warmup == WarmupMode::None {
189            snapshot
190        } else {
191            validate_protocol(snapshot.protocol())?;
192            materialize_eager_scan_metadata(snapshot).await?
193        };
194        Ok(DeltaTable::new(snapshot, self.execution_options))
195    }
196
197    /// Loads a Delta Kernel snapshot without converting its logical Arrow schema.
198    ///
199    /// Warmup settings apply to [`DeltaTable`] and are not supported by this method. Calling this
200    /// method after selecting any mode other than [`WarmupMode::None`] returns a configuration
201    /// error.
202    ///
203    /// # Errors
204    ///
205    /// Returns an error if warmup was requested or if the table location, storage, or snapshot
206    /// cannot be loaded.
207    pub async fn load_snapshot(self) -> Result<DeltaTableSnapshot, DeltaReaderError> {
208        if self.warmup != WarmupMode::None {
209            return InvalidConfigurationSnafu {
210                reason: "snapshot_load_does_not_support_table_warmup",
211            }
212            .fail();
213        }
214        let snapshot = load_kernel_table_snapshot(
215            self.table_location,
216            self.storage_options,
217            self.snapshot_selection,
218        )
219        .await?;
220        Ok(DeltaTableSnapshot::new(snapshot, self.execution_options))
221    }
222}
223
224async fn materialize_eager_scan_metadata(
225    snapshot: ArrowTableSnapshot,
226) -> Result<ArrowTableSnapshot, DeltaReaderError> {
227    tokio::task::spawn_blocking(move || snapshot.materialize_eager_scan_metadata())
228        .await
229        .boxed()
230        .context(ScanPlanningSnafu {
231            reason: "query_planning_warmup_task_failed",
232        })?
233}
234
235impl fmt::Debug for DeltaTableBuilder {
236    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
237        formatter
238            .debug_struct("DeltaTableBuilder")
239            .field("table_location", &"<redacted>")
240            .field("storage_options", &"<redacted>")
241            .field("snapshot_selection", &self.snapshot_selection)
242            .field("execution_options", &self.execution_options)
243            .field("warmup", &self.warmup)
244            .finish()
245    }
246}
247
248/// Loaded Delta snapshot metadata awaiting logical Arrow schema conversion.
249pub struct DeltaTableSnapshot {
250    snapshot: KernelTableSnapshot,
251    execution_options: DeltaScanExecutionOptions,
252}
253
254impl DeltaTableSnapshot {
255    fn new(snapshot: KernelTableSnapshot, execution_options: DeltaScanExecutionOptions) -> Self {
256        Self {
257            snapshot,
258            execution_options,
259        }
260    }
261
262    /// Returns the loaded Delta snapshot version.
263    pub fn version(&self) -> u64 {
264        self.snapshot.version()
265    }
266
267    /// Returns the loaded Delta protocol metadata.
268    pub fn protocol(&self) -> &DeltaProtocol {
269        self.snapshot.protocol()
270    }
271
272    /// Returns the normalized table URL.
273    ///
274    /// This value may contain sensitive caller input. Do not log or expose it.
275    pub fn table_url(&self) -> &str {
276        self.snapshot.table_url()
277    }
278
279    /// Validates the loaded snapshot against the supported reader protocol.
280    pub fn validate_protocol(&self) -> Result<(), DeltaReaderError> {
281        validate_protocol(self.protocol())
282    }
283
284    /// Converts the logical Arrow schema and finishes constructing the table.
285    pub fn into_table(self) -> Result<DeltaTable, DeltaReaderError> {
286        Ok(DeltaTable::new(
287            self.snapshot.into_arrow_snapshot()?,
288            self.execution_options,
289        ))
290    }
291}
292
293impl fmt::Debug for DeltaTableSnapshot {
294    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
295        formatter
296            .debug_struct("DeltaTableSnapshot")
297            .field("version", &self.version())
298            .finish_non_exhaustive()
299    }
300}
301
302/// One immutable loaded Delta table snapshot.
303#[derive(Clone)]
304pub struct DeltaTable {
305    snapshot: Arc<ArrowTableSnapshot>,
306    execution_options: DeltaScanExecutionOptions,
307}
308
309impl DeltaTable {
310    fn new(snapshot: ArrowTableSnapshot, execution_options: DeltaScanExecutionOptions) -> Self {
311        Self {
312            snapshot: Arc::new(snapshot),
313            execution_options,
314        }
315    }
316
317    /// Returns the loaded Delta snapshot version.
318    pub fn version(&self) -> u64 {
319        self.snapshot.version()
320    }
321
322    /// Returns a shared handle to the logical Arrow schema.
323    pub fn schema(&self) -> SchemaRef {
324        self.snapshot.schema()
325    }
326
327    /// Returns the loaded Delta protocol metadata.
328    pub fn protocol(&self) -> &DeltaProtocol {
329        self.snapshot.protocol()
330    }
331
332    /// Returns the normalized table URL.
333    ///
334    /// This value may contain sensitive caller input. Do not log or expose it.
335    pub fn table_url(&self) -> &str {
336        self.snapshot.table_url()
337    }
338
339    #[allow(dead_code)]
340    pub(crate) fn partition_columns(&self) -> &[String] {
341        self.snapshot.partition_columns()
342    }
343
344    #[allow(dead_code)]
345    pub(crate) fn snapshot(&self) -> &ArrowTableSnapshot {
346        self.snapshot.as_ref()
347    }
348
349    /// Validates the loaded snapshot against the supported reader protocol.
350    pub fn validate_protocol(&self) -> Result<(), DeltaReaderError> {
351        validate_protocol(self.protocol())
352    }
353
354    /// Loads the latest version from this table and returns it as a new immutable table.
355    ///
356    /// The current table and any scans built from it remain fixed at their original version. A
357    /// query-planning metadata cache is refreshed for the new version and reused directly when the
358    /// version has not changed.
359    ///
360    /// # Errors
361    ///
362    /// Returns an error if the newer snapshot, schema, or retained query-planning metadata cannot
363    /// be loaded. No partially refreshed table is returned.
364    pub async fn refresh(&self) -> Result<Self, DeltaReaderError> {
365        let snapshot = Arc::clone(&self.snapshot);
366        let snapshot = tokio::task::spawn_blocking(move || snapshot.refresh())
367            .await
368            .boxed()
369            .context(SnapshotLoadSnafu {
370                reason: "snapshot_refresh_task_failed",
371            })
372            .and_then(|result| result)?;
373        Ok(Self::new(snapshot, self.execution_options))
374    }
375
376    /// Starts configuring a new single-use scan.
377    pub fn scan(&self) -> DeltaScanBuilder<'_> {
378        DeltaScanBuilder {
379            table: self,
380            projection: None,
381            predicate: None,
382            limit: None,
383            target_partitions: None,
384            execution_options: self.execution_options,
385        }
386    }
387}
388
389impl fmt::Debug for DeltaTable {
390    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
391        formatter
392            .debug_struct("DeltaTable")
393            .field("version", &self.version())
394            .finish_non_exhaustive()
395    }
396}
397
398/// Configures one single-use streaming Delta scan.
399#[must_use = "scan builder settings do nothing unless the scan is built"]
400pub struct DeltaScanBuilder<'table> {
401    table: &'table DeltaTable,
402    projection: Option<Vec<String>>,
403    predicate: Option<DeltaPredicate>,
404    limit: Option<usize>,
405    target_partitions: Option<usize>,
406    execution_options: DeltaScanExecutionOptions,
407}
408
409impl<'table> DeltaScanBuilder<'table> {
410    /// Selects visible logical columns in caller order.
411    pub fn with_projection(
412        mut self,
413        logical_columns: impl IntoIterator<Item = impl Into<String>>,
414    ) -> Self {
415        self.projection = Some(logical_columns.into_iter().map(Into::into).collect());
416        self
417    }
418
419    /// Replaces the exact logical row predicate.
420    pub fn with_predicate(mut self, predicate: DeltaPredicate) -> Self {
421        self.predicate = Some(predicate);
422        self
423    }
424
425    /// Sets the maximum number of output rows.
426    pub const fn with_limit(mut self, limit: usize) -> Self {
427        self.limit = Some(limit);
428        self
429    }
430
431    /// Overrides the number of planned scan partitions.
432    pub fn with_target_partitions(
433        mut self,
434        target_partitions: usize,
435    ) -> Result<Self, DeltaReaderError> {
436        if target_partitions == 0 {
437            return InvalidConfigurationSnafu {
438                reason: "scan_partition_target_must_be_positive",
439            }
440            .fail();
441        }
442        self.target_partitions = Some(target_partitions);
443        Ok(self)
444    }
445
446    /// Replaces the execution settings for this scan.
447    pub const fn with_execution_options(
448        mut self,
449        execution_options: DeltaScanExecutionOptions,
450    ) -> Self {
451        self.execution_options = execution_options;
452        self
453    }
454
455    /// Builds one immutable single-use scan plan without reading data files.
456    pub async fn build(self) -> Result<DeltaScan, DeltaReaderError> {
457        self.table.validate_protocol()?;
458        if let Some(predicate) = self.predicate.as_ref() {
459            validate_predicate(predicate, self.table.schema().as_ref())?;
460        }
461
462        let snapshot_version = self.table.version();
463        let backend = self.execution_options.parquet_backend();
464        let scan_metadata_source = if self.table.snapshot.eager_scan_metadata().is_some() {
465            QUERY_PLANNING_CACHE_SCAN_METADATA_SOURCE
466        } else {
467            DELTA_LOG_SCAN_METADATA_SOURCE
468        };
469        trace_planning_started(snapshot_version, backend, scan_metadata_source);
470        let snapshot = Arc::clone(&self.table.snapshot);
471        let projection = self.projection;
472        let predicate = self.predicate;
473        let hidden_columns = predicate
474            .as_ref()
475            .map(referenced_columns)
476            .unwrap_or_default();
477        let enforce_physical_predicate_rows =
478            predicate.as_ref().is_some_and(kernel_pruning_is_exact);
479        let kernel_predicate = predicate.as_ref().and_then(kernel_pruning_predicate);
480        let include_stats = kernel_predicate.is_some();
481        let execution_options = self.execution_options;
482        let target_partitions = self.target_partitions;
483        let result = tokio::task::spawn_blocking(move || {
484            plan_scan(
485                snapshot.as_ref(),
486                projection.as_deref(),
487                &hidden_columns,
488                kernel_predicate,
489                include_stats,
490                execution_options,
491                DeltaScanPartitionTargetOptions {
492                    explicit_target_partitions: target_partitions,
493                    datafusion_target_partitions: None,
494                },
495            )
496        })
497        .await
498        .boxed()
499        .context(ScanPlanningSnafu {
500            reason: "scan_planning_task_failed",
501        })
502        .and_then(|result| result);
503
504        match result {
505            Ok(plan) => {
506                trace_planning_completed(
507                    snapshot_version,
508                    backend,
509                    plan.partitions.len(),
510                    scan_metadata_source,
511                );
512                Ok(DeltaScan {
513                    plan: Arc::new(plan),
514                    predicate,
515                    limit: self.limit,
516                    enforce_physical_predicate_rows,
517                })
518            }
519            Err(error) => {
520                trace_planning_failed(snapshot_version, backend, scan_metadata_source, &error);
521                Err(error)
522            }
523        }
524    }
525}
526
527/// One immutable, single-use streaming Delta scan plan.
528///
529/// A scan cannot be cloned or converted into a stream twice.
530///
531/// ```compile_fail
532/// use delta_arrow_reader::DeltaScan;
533///
534/// fn stream_twice(scan: DeltaScan) {
535///     let _ = scan.into_stream();
536///     let _ = scan.into_stream();
537/// }
538/// ```
539///
540/// ```compile_fail
541/// use delta_arrow_reader::DeltaScan;
542///
543/// fn clone_scan(scan: DeltaScan) {
544///     let _ = scan.clone();
545/// }
546/// ```
547#[must_use = "scans do nothing unless converted into a stream"]
548pub struct DeltaScan {
549    plan: Arc<DeltaScanPlan>,
550    predicate: Option<DeltaPredicate>,
551    limit: Option<usize>,
552    enforce_physical_predicate_rows: bool,
553}
554
555impl DeltaScan {
556    /// Returns a shared handle to the visible logical output schema.
557    pub fn schema(&self) -> SchemaRef {
558        Arc::clone(&self.plan.projected_schema)
559    }
560
561    /// Returns the number of planned execution partitions.
562    pub fn partition_count(&self) -> usize {
563        self.plan.partitions.len()
564    }
565
566    /// Converts the scan into a pull-driven Arrow batch stream.
567    ///
568    /// Data-file reads begin only when the stream is polled.
569    pub fn into_stream(self) -> DeltaBatchStream {
570        let metrics = self.plan.metrics.clone();
571        let schema = Arc::clone(&self.plan.projected_schema);
572        let partition_count = self.plan.partitions.len();
573        let snapshot_version = self.plan.snapshot_version;
574        let backend = self.plan.execution_options.parquet_backend();
575        let projection = (self.plan.logical_schema.as_ref() != schema.as_ref())
576            .then(|| (0..schema.fields().len()).collect::<Vec<_>>());
577        let partitions = if self.limit == Some(0) {
578            VecDeque::new()
579        } else {
580            let scheduler = DeltaScanScheduler::new(Arc::clone(&self.plan));
581            let admission: FileAdmissionPolicy<_> = Arc::new(|_| Ok(FileAdmissionDecision::Admit));
582            let executor = match backend {
583                ParquetReaderBackend::Direct => direct_parquet_executor(
584                    &self.plan,
585                    None,
586                    self.enforce_physical_predicate_rows
587                        .then(|| self.plan.physical_predicate.clone())
588                        .flatten(),
589                ),
590                ParquetReaderBackend::DeltaKernel => delta_kernel_executor(&self.plan),
591            };
592            scheduler.partition_streams(admission, executor)
593        };
594
595        DeltaBatchStream {
596            schema,
597            metrics,
598            partitions,
599            predicate: self.predicate,
600            projection,
601            remaining: self.limit,
602            snapshot_version,
603            backend,
604            partition_count,
605            started: false,
606            done: false,
607        }
608    }
609}
610
611/// Pull-driven stream of finalized logical Arrow batches from one Delta scan.
612///
613/// The stream has no inherent whole-result collection method. Callers that
614/// intentionally materialize a result must opt into a stream extension trait.
615///
616/// ```compile_fail
617/// use delta_arrow_reader::DeltaBatchStream;
618///
619/// fn collect_without_opt_in(stream: DeltaBatchStream) {
620///     let _ = stream.collect();
621/// }
622/// ```
623#[must_use = "streams do nothing unless polled"]
624pub struct DeltaBatchStream {
625    schema: SchemaRef,
626    metrics: DeltaScanMetrics,
627    partitions: VecDeque<PartitionStream>,
628    predicate: Option<DeltaPredicate>,
629    projection: Option<Vec<usize>>,
630    remaining: Option<usize>,
631    snapshot_version: u64,
632    backend: ParquetReaderBackend,
633    partition_count: usize,
634    started: bool,
635    done: bool,
636}
637
638impl DeltaBatchStream {
639    /// Returns a shared handle to the visible logical output schema.
640    pub fn schema(&self) -> SchemaRef {
641        Arc::clone(&self.schema)
642    }
643
644    /// Returns a lightweight shared handle to live scan metrics.
645    pub fn metrics(&self) -> DeltaScanMetrics {
646        self.metrics.clone()
647    }
648
649    fn start(&mut self) {
650        if self.started {
651            return;
652        }
653        self.started = true;
654        trace_execution_started(self.snapshot_version, self.backend, self.partition_count);
655        for partition in &mut self.partitions {
656            partition.start();
657        }
658    }
659
660    fn complete(&mut self) {
661        if self.done {
662            return;
663        }
664        self.partitions.clear();
665        self.done = true;
666        trace_execution_completed(self.snapshot_version, self.backend, self.partition_count);
667    }
668
669    fn fail(&mut self, error: &DeltaReaderError) {
670        self.partitions.clear();
671        self.done = true;
672        trace_execution_failed(
673            self.snapshot_version,
674            self.backend,
675            self.partition_count,
676            error,
677        );
678    }
679
680    fn finalize_batch(&self, mut batch: RecordBatch) -> Result<RecordBatch, DeltaReaderError> {
681        if let Some(predicate) = self.predicate.as_ref() {
682            batch = evaluate_predicate(&batch, predicate)?;
683        }
684        if let Some(projection) = self.projection.as_ref() {
685            batch = batch
686                .project(projection)
687                .boxed()
688                .context(DataFileReadSnafu {
689                    reason: "direct_projection_failed",
690                })?;
691        }
692        Ok(batch)
693    }
694}
695
696impl Stream for DeltaBatchStream {
697    type Item = Result<RecordBatch, DeltaReaderError>;
698
699    fn poll_next(self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<Option<Self::Item>> {
700        let this = self.get_mut();
701        if this.done {
702            return Poll::Ready(None);
703        }
704        this.start();
705
706        loop {
707            let Some(partition) = this.partitions.front_mut() else {
708                this.complete();
709                return Poll::Ready(None);
710            };
711            match Pin::new(partition).poll_next(context) {
712                Poll::Ready(Some(Ok(batch))) => {
713                    let mut batch = match this.finalize_batch(batch) {
714                        Ok(batch) => batch,
715                        Err(error) => {
716                            this.fail(&error);
717                            return Poll::Ready(Some(Err(error)));
718                        }
719                    };
720                    if let Some(remaining) = this.remaining.as_mut() {
721                        if batch.num_rows() >= *remaining {
722                            batch = batch.slice(0, *remaining);
723                            *remaining = 0;
724                            this.complete();
725                        } else {
726                            *remaining -= batch.num_rows();
727                        }
728                    }
729                    return Poll::Ready(Some(Ok(batch)));
730                }
731                Poll::Ready(Some(Err(error))) => {
732                    this.fail(&error);
733                    return Poll::Ready(Some(Err(error)));
734                }
735                Poll::Ready(None) => {
736                    this.partitions.pop_front();
737                }
738                Poll::Pending => return Poll::Pending,
739            }
740        }
741    }
742}
743
744impl Drop for DeltaBatchStream {
745    fn drop(&mut self) {
746        if self.done {
747            return;
748        }
749        self.partitions.clear();
750        self.done = true;
751        trace_execution_dropped(self.snapshot_version, self.backend, self.partition_count);
752    }
753}
754
755pub(crate) fn direct_parquet_executor(
756    plan: &Arc<DeltaScanPlan>,
757    output_batch_size_rows: Option<usize>,
758    row_predicate: Option<crate::delta::kernel::DeltaKernelPredicate>,
759) -> FileExecutor<planning::DeltaScanFileTask, FileBatchStream> {
760    backend::direct_parquet::direct_parquet_file_executor(
761        plan,
762        output_batch_size_rows,
763        row_predicate,
764        Arc::default(),
765        None,
766    )
767}
768
769pub(crate) fn delta_kernel_executor(
770    plan: &Arc<DeltaScanPlan>,
771) -> FileExecutor<planning::DeltaScanFileTask, FileBatchStream> {
772    backend::kernel_reader::delta_kernel_file_executor(plan)
773}
774
775fn trace_planning_started(
776    snapshot_version: u64,
777    backend: ParquetReaderBackend,
778    scan_metadata_source: &'static str,
779) {
780    tracing::debug!(
781        target: TRACING_TARGET,
782        event = "scan_planning.started",
783        snapshot_version,
784        backend = ?backend,
785        partition_count = tracing::field::Empty,
786        scan_metadata_source,
787        outcome = "started"
788    );
789}
790
791fn trace_planning_completed(
792    snapshot_version: u64,
793    backend: ParquetReaderBackend,
794    partition_count: usize,
795    scan_metadata_source: &'static str,
796) {
797    tracing::debug!(
798        target: TRACING_TARGET,
799        event = "scan_planning.completed",
800        snapshot_version,
801        backend = ?backend,
802        partition_count,
803        scan_metadata_source,
804        outcome = "completed"
805    );
806}
807
808fn trace_planning_failed(
809    snapshot_version: u64,
810    backend: ParquetReaderBackend,
811    scan_metadata_source: &'static str,
812    error: &DeltaReaderError,
813) {
814    tracing::debug!(
815        target: TRACING_TARGET,
816        event = "scan_planning.failed",
817        snapshot_version,
818        backend = ?backend,
819        partition_count = tracing::field::Empty,
820        scan_metadata_source,
821        outcome = "failed",
822        error_code = error.code(),
823        error_phase = error.phase().as_str()
824    );
825}
826
827fn trace_execution_started(
828    snapshot_version: u64,
829    backend: ParquetReaderBackend,
830    partition_count: usize,
831) {
832    tracing::debug!(
833        target: TRACING_TARGET,
834        event = "scan_execution.started",
835        snapshot_version,
836        backend = ?backend,
837        partition_count,
838        outcome = "started"
839    );
840}
841
842fn trace_execution_completed(
843    snapshot_version: u64,
844    backend: ParquetReaderBackend,
845    partition_count: usize,
846) {
847    tracing::debug!(
848        target: TRACING_TARGET,
849        event = "scan_execution.completed",
850        snapshot_version,
851        backend = ?backend,
852        partition_count,
853        outcome = "completed"
854    );
855}
856
857fn trace_execution_failed(
858    snapshot_version: u64,
859    backend: ParquetReaderBackend,
860    partition_count: usize,
861    error: &DeltaReaderError,
862) {
863    tracing::debug!(
864        target: TRACING_TARGET,
865        event = "scan_execution.failed",
866        snapshot_version,
867        backend = ?backend,
868        partition_count,
869        outcome = "failed",
870        error_code = error.code(),
871        error_phase = error.phase().as_str()
872    );
873}
874
875fn trace_execution_dropped(
876    snapshot_version: u64,
877    backend: ParquetReaderBackend,
878    partition_count: usize,
879) {
880    tracing::debug!(
881        target: TRACING_TARGET,
882        event = "scan_execution.dropped",
883        snapshot_version,
884        backend = ?backend,
885        partition_count,
886        outcome = "dropped"
887    );
888}
889
890#[cfg(test)]
891mod tests {
892    use std::{
893        collections::{BTreeMap, VecDeque},
894        fmt, fs,
895        future::pending,
896        path::{Path, PathBuf},
897        sync::{Arc, Mutex, Once},
898        time::{Duration, SystemTime, UNIX_EPOCH},
899    };
900
901    use arrow::{
902        array::Int32Array,
903        datatypes::{DataType, Field, Schema, SchemaRef},
904        record_batch::RecordBatch,
905    };
906    use futures_util::{FutureExt, StreamExt, stream};
907    use tokio::{sync::Notify, time::timeout};
908    use tracing::{
909        Event, Level, Metadata, Subscriber,
910        field::{Field as TracingField, Visit},
911        span::{Attributes, Id, Record},
912        subscriber::{Interest, with_default},
913    };
914
915    use super::{
916        DeltaBatchStream, DeltaTable, trace_execution_completed, trace_execution_dropped,
917        trace_execution_failed, trace_execution_started, trace_planning_completed,
918        trace_planning_failed, trace_planning_started,
919    };
920    use crate::{
921        DeltaScanExecutionOptions, DeltaScanMetrics, DeltaSnapshotSelection, DeltaStorageOptions,
922        ParquetReaderBackend,
923        delta::snapshot::load_delta_table_snapshot_blocking,
924        error::InvalidConfigurationSnafu,
925        reader::{
926            metrics::DeltaScanMetricsConfig,
927            scheduling::{
928                FileAdmissionDecision, FileAdmissionPolicy, FileBatchStream, FileExecutor,
929                FileReadPermit, PartitionStream, ScanCancellation, ScanReadLimiter,
930            },
931        },
932    };
933
934    static TRACING_TEST_LOCK: Mutex<()> = Mutex::new(());
935    static TRACING_TEST_GLOBAL_SUBSCRIBER: Once = Once::new();
936
937    #[derive(Clone, Default)]
938    struct EventFields(Arc<Mutex<Vec<BTreeMap<String, String>>>>);
939
940    impl Subscriber for EventFields {
941        fn register_callsite(&self, metadata: &'static Metadata<'static>) -> Interest {
942            if metadata.target() == "delta_arrow_reader" && *metadata.level() == Level::DEBUG {
943                Interest::always()
944            } else {
945                Interest::sometimes()
946            }
947        }
948
949        fn enabled(&self, metadata: &Metadata<'_>) -> bool {
950            metadata.target() == "delta_arrow_reader" && *metadata.level() == Level::DEBUG
951        }
952
953        fn new_span(&self, _attributes: &Attributes<'_>) -> Id {
954            Id::from_u64(1)
955        }
956
957        fn record(&self, _span: &Id, _values: &Record<'_>) {}
958
959        fn record_follows_from(&self, _span: &Id, _follows: &Id) {}
960
961        fn event(&self, event: &Event<'_>) {
962            let metadata = event.metadata();
963            assert_eq!(metadata.target(), "delta_arrow_reader");
964            let mut fields = metadata
965                .fields()
966                .iter()
967                .map(|field| (field.name().to_owned(), "<empty>".to_owned()))
968                .collect();
969            event.record(&mut FieldVisitor(&mut fields));
970            self.0.lock().expect("event lock").push(fields);
971        }
972
973        fn enter(&self, _span: &Id) {}
974
975        fn exit(&self, _span: &Id) {}
976    }
977
978    struct FieldVisitor<'fields>(&'fields mut BTreeMap<String, String>);
979
980    impl Visit for FieldVisitor<'_> {
981        fn record_debug(&mut self, field: &TracingField, value: &dyn fmt::Debug) {
982            self.0.insert(field.name().to_owned(), format!("{value:?}"));
983        }
984
985        fn record_str(&mut self, field: &TracingField, value: &str) {
986            self.0.insert(field.name().to_owned(), value.to_owned());
987        }
988
989        fn record_u64(&mut self, field: &TracingField, value: u64) {
990            self.0.insert(field.name().to_owned(), value.to_string());
991        }
992
993        fn record_i64(&mut self, field: &TracingField, value: i64) {
994            self.0.insert(field.name().to_owned(), value.to_string());
995        }
996
997        fn record_u128(&mut self, field: &TracingField, value: u128) {
998            self.0.insert(field.name().to_owned(), value.to_string());
999        }
1000    }
1001
1002    struct DeltaLogTable(PathBuf);
1003
1004    impl DeltaLogTable {
1005        fn new(name: &str) -> Result<Self, Box<dyn std::error::Error>> {
1006            let nanos = SystemTime::now().duration_since(UNIX_EPOCH)?.as_nanos();
1007            let path = Path::new("target")
1008                .join("delta-arrow-reader-tracing-tests")
1009                .join(format!("{}-{name}-{nanos}", std::process::id()));
1010            fs::create_dir_all(path.join("_delta_log"))?;
1011            fs::write(
1012                path.join("_delta_log/00000000000000000000.json"),
1013                r#"{"protocol":{"minReaderVersion":1,"minWriterVersion":2}}
1014{"metaData":{"id":"tracing-test","format":{"provider":"parquet","options":{}},"schemaString":"{\"type\":\"struct\",\"fields\":[{\"name\":\"id\",\"type\":\"integer\",\"nullable\":true,\"metadata\":{}}]}","partitionColumns":[],"configuration":{},"createdTime":1587968585495}}
1015{"add":{"path":"secret-planning-object.parquet","partitionValues":{},"size":10,"modificationTime":1587968586000,"dataChange":true}}
1016"#,
1017            )?;
1018            Ok(Self(path))
1019        }
1020    }
1021
1022    impl Drop for DeltaLogTable {
1023        fn drop(&mut self) {
1024            let _ = fs::remove_dir_all(&self.0);
1025        }
1026    }
1027
1028    fn capture_events<T>(run: impl FnOnce() -> T) -> (T, Vec<BTreeMap<String, String>>) {
1029        let _lock = TRACING_TEST_LOCK
1030            .lock()
1031            .unwrap_or_else(std::sync::PoisonError::into_inner);
1032        let events = Arc::new(Mutex::new(Vec::new()));
1033        let subscriber = EventFields(Arc::clone(&events));
1034        TRACING_TEST_GLOBAL_SUBSCRIBER.call_once(|| {
1035            let _ = tracing::subscriber::set_global_default(EventFields::default());
1036        });
1037        let result = with_default(subscriber, || {
1038            tracing::callsite::rebuild_interest_cache();
1039            run()
1040        });
1041        tracing::callsite::rebuild_interest_cache();
1042        let captured = events
1043            .lock()
1044            .map(|events| events.clone())
1045            .unwrap_or_default();
1046        (result, captured)
1047    }
1048
1049    struct ControlledMerge {
1050        stream: DeltaBatchStream,
1051        limiter: Arc<ScanReadLimiter>,
1052        cancellation: ScanCancellation,
1053        metrics: DeltaScanMetrics,
1054        first_partition_gate: Arc<Notify>,
1055    }
1056
1057    fn schema() -> SchemaRef {
1058        Arc::new(Schema::new(vec![Field::new("id", DataType::Int32, false)]))
1059    }
1060
1061    fn batch(id: i32) -> RecordBatch {
1062        RecordBatch::try_new(schema(), vec![Arc::new(Int32Array::from(vec![id]))])
1063            .expect("valid test batch")
1064    }
1065
1066    fn batch_id(batch: &RecordBatch) -> i32 {
1067        batch
1068            .column(0)
1069            .as_any()
1070            .downcast_ref::<Int32Array>()
1071            .expect("Int32 id")
1072            .value(0)
1073    }
1074
1075    fn execution_options() -> Result<DeltaScanExecutionOptions, crate::DeltaReaderError> {
1076        DeltaScanExecutionOptions::new()
1077            .with_prefetch_files_per_partition(0)
1078            .with_max_concurrent_file_reads_per_partition(1)?
1079            .with_max_concurrent_file_reads_per_scan(Some(2))?
1080            .with_output_buffer_batches_per_partition(1)
1081    }
1082
1083    fn metrics() -> DeltaScanMetrics {
1084        DeltaScanMetrics::new(DeltaScanMetricsConfig {
1085            snapshot_version: 7,
1086            parquet_backend: ParquetReaderBackend::Direct,
1087            scan_partitions_planned: 2,
1088            files_planned: 2,
1089            add_actions_excluded_during_planning: Some(0),
1090            estimated_input_rows: Some(4),
1091            estimated_input_bytes: Some(4),
1092        })
1093    }
1094
1095    fn file_stream(permit: FileReadPermit, batches: Vec<RecordBatch>) -> FileBatchStream {
1096        Box::pin(stream::unfold(
1097            (VecDeque::from(batches), permit),
1098            |(mut batches, permit)| async move {
1099                batches
1100                    .pop_front()
1101                    .map(|batch| (Ok(batch), (batches, permit)))
1102            },
1103        ))
1104    }
1105
1106    fn gated_file_stream(
1107        permit: FileReadPermit,
1108        batches: Vec<RecordBatch>,
1109        gate: Arc<Notify>,
1110    ) -> FileBatchStream {
1111        Box::pin(stream::unfold(
1112            (false, VecDeque::from(batches), permit, gate),
1113            |(wait, mut batches, permit, gate)| async move {
1114                let batch = batches.pop_front()?;
1115                if wait {
1116                    gate.notified().await;
1117                }
1118                Some((Ok(batch), (true, batches, permit, gate)))
1119            },
1120        ))
1121    }
1122
1123    fn direct_stream(
1124        partitions: VecDeque<PartitionStream>,
1125        metrics: DeltaScanMetrics,
1126    ) -> DeltaBatchStream {
1127        DeltaBatchStream {
1128            schema: schema(),
1129            metrics,
1130            partitions,
1131            predicate: None,
1132            projection: None,
1133            remaining: None,
1134            snapshot_version: 7,
1135            backend: ParquetReaderBackend::Direct,
1136            partition_count: 2,
1137            started: false,
1138            done: false,
1139        }
1140    }
1141
1142    fn controlled_merge() -> Result<ControlledMerge, Box<dyn std::error::Error>> {
1143        let options = execution_options()?;
1144        let limiter = ScanReadLimiter::new(options, 2, 2);
1145        let cancellation = ScanCancellation::new();
1146        let metrics = metrics();
1147        let first_partition_gate = Arc::new(Notify::new());
1148        let executor: FileExecutor<i32, FileBatchStream> = {
1149            let gate = Arc::clone(&first_partition_gate);
1150            Arc::new(move |task, permit, _| {
1151                let gate = Arc::clone(&gate);
1152                async move {
1153                    let batches = vec![batch(task), batch(task * 2)];
1154                    Ok(if task == 1 {
1155                        gated_file_stream(permit, batches, gate)
1156                    } else {
1157                        file_stream(permit, batches)
1158                    })
1159                }
1160                .boxed()
1161            })
1162        };
1163        let admission: FileAdmissionPolicy<i32> =
1164            Arc::new(|_: &i32| Ok(FileAdmissionDecision::Admit));
1165        let first = PartitionStream::new(
1166            vec![1],
1167            limiter.partition(0)?,
1168            options,
1169            admission.clone(),
1170            Arc::clone(&executor),
1171            metrics.clone(),
1172            cancellation.clone(),
1173        );
1174        let second = PartitionStream::new(
1175            vec![10],
1176            limiter.partition(1)?,
1177            options,
1178            admission,
1179            executor,
1180            metrics.clone(),
1181            cancellation.clone(),
1182        );
1183
1184        Ok(ControlledMerge {
1185            stream: direct_stream(VecDeque::from([first, second]), metrics.clone()),
1186            limiter,
1187            cancellation,
1188            metrics,
1189            first_partition_gate,
1190        })
1191    }
1192
1193    async fn wait_for_batches(metrics: &DeltaScanMetrics, expected: u64) {
1194        timeout(Duration::from_secs(5), async {
1195            while metrics.snapshot().scheduler_batches_emitted < expected {
1196                tokio::task::yield_now().await;
1197            }
1198        })
1199        .await
1200        .expect("batch production reached expected bound");
1201    }
1202
1203    #[test]
1204    fn lifecycle_tracing_has_only_bounded_fields() {
1205        let error = InvalidConfigurationSnafu { reason: "test" }.build();
1206
1207        let (_, events) = capture_events(|| {
1208            trace_planning_started(7, ParquetReaderBackend::Direct, "query_planning_cache");
1209            trace_planning_completed(7, ParquetReaderBackend::Direct, 2, "query_planning_cache");
1210            trace_planning_failed(
1211                7,
1212                ParquetReaderBackend::Direct,
1213                "query_planning_cache",
1214                &error,
1215            );
1216            trace_execution_started(7, ParquetReaderBackend::Direct, 2);
1217            trace_execution_completed(7, ParquetReaderBackend::Direct, 2);
1218            trace_execution_failed(7, ParquetReaderBackend::Direct, 2, &error);
1219            trace_execution_dropped(7, ParquetReaderBackend::Direct, 2);
1220        });
1221
1222        assert_eq!(events.len(), 7);
1223        let allowed = [
1224            "backend",
1225            "error_phase",
1226            "error_code",
1227            "event",
1228            "outcome",
1229            "partition_count",
1230            "scan_metadata_source",
1231            "snapshot_version",
1232        ];
1233        for fields in events.iter() {
1234            assert!(fields.keys().all(|field| allowed.contains(&field.as_str())));
1235            assert!(fields.contains_key("event"));
1236            assert!(fields.contains_key("snapshot_version"));
1237            assert!(fields.contains_key("backend"));
1238            assert!(fields.contains_key("partition_count"));
1239            assert!(fields.contains_key("outcome"));
1240            if fields
1241                .get("event")
1242                .is_some_and(|event| event.starts_with("scan_planning."))
1243            {
1244                assert_eq!(
1245                    fields.get("scan_metadata_source").map(String::as_str),
1246                    Some("query_planning_cache")
1247                );
1248            }
1249        }
1250    }
1251
1252    #[test]
1253    fn planning_tracing_reports_the_table_metadata_source_on_success_and_failure()
1254    -> Result<(), Box<dyn std::error::Error>> {
1255        const OBJECT_KEY: &str = "secret-planning-object.parquet";
1256        const STORAGE_VALUE: &str = "secret-planning-storage-value";
1257        let fixture = DeltaLogTable::new("planning-source")?;
1258        let mut storage_options = DeltaStorageOptions::new();
1259        storage_options.insert("secret-option".to_owned(), STORAGE_VALUE.to_owned());
1260        let snapshot = load_delta_table_snapshot_blocking(
1261            &fixture.0.to_string_lossy(),
1262            &storage_options,
1263            DeltaSnapshotSelection::Latest,
1264        )?;
1265        let eager_snapshot = snapshot.clone().materialize_eager_scan_metadata()?;
1266        let lazy = DeltaTable::new(snapshot, DeltaScanExecutionOptions::new());
1267        let eager = DeltaTable::new(eager_snapshot, DeltaScanExecutionOptions::new());
1268        let runtime = tokio::runtime::Builder::new_current_thread()
1269            .enable_all()
1270            .build()?;
1271
1272        let (eager_result, eager_events) =
1273            capture_events(|| runtime.block_on(eager.scan().build()));
1274        let _ = eager_result?;
1275        assert_eq!(eager_events.len(), 2);
1276        assert_eq!(
1277            eager_events[0].get("event").map(String::as_str),
1278            Some("scan_planning.started")
1279        );
1280        assert_eq!(
1281            eager_events[1].get("event").map(String::as_str),
1282            Some("scan_planning.completed")
1283        );
1284        assert!(eager_events.iter().all(|event| {
1285            event.get("scan_metadata_source").map(String::as_str) == Some("query_planning_cache")
1286        }));
1287
1288        let (failed_result, failed_events) =
1289            capture_events(|| runtime.block_on(eager.scan().with_projection(["missing"]).build()));
1290        assert!(failed_result.is_err());
1291        assert_eq!(failed_events.len(), 2);
1292        assert_eq!(
1293            failed_events[0].get("event").map(String::as_str),
1294            Some("scan_planning.started")
1295        );
1296        assert_eq!(
1297            failed_events[1].get("event").map(String::as_str),
1298            Some("scan_planning.failed")
1299        );
1300        assert!(failed_events.iter().all(|event| {
1301            event.get("scan_metadata_source").map(String::as_str) == Some("query_planning_cache")
1302        }));
1303
1304        let (lazy_result, lazy_events) = capture_events(|| runtime.block_on(lazy.scan().build()));
1305        let _ = lazy_result?;
1306        assert_eq!(lazy_events.len(), 2);
1307        assert!(lazy_events.iter().all(|event| {
1308            event.get("scan_metadata_source").map(String::as_str) == Some("delta_log")
1309        }));
1310
1311        let captured = format!("{eager_events:?}{failed_events:?}{lazy_events:?}");
1312        assert!(!captured.contains(&fixture.0.to_string_lossy().into_owned()));
1313        assert!(!captured.contains(OBJECT_KEY));
1314        assert!(!captured.contains(STORAGE_VALUE));
1315        Ok(())
1316    }
1317
1318    #[tokio::test]
1319    async fn merged_stream_is_ordered_and_bounds_later_partition_queues()
1320    -> Result<(), Box<dyn std::error::Error>> {
1321        let ControlledMerge {
1322            mut stream,
1323            limiter,
1324            metrics,
1325            first_partition_gate,
1326            ..
1327        } = controlled_merge()?;
1328
1329        let first = stream.next().await.ok_or("first batch missing")??;
1330        assert_eq!(batch_id(&first), 1);
1331        wait_for_batches(&metrics, 2).await;
1332        for _ in 0..32 {
1333            tokio::task::yield_now().await;
1334        }
1335        assert_eq!(metrics.snapshot().scheduler_batches_emitted, 2);
1336        assert_eq!(limiter.active_file_reads(), 2);
1337
1338        first_partition_gate.notify_one();
1339        let mut ids = vec![batch_id(
1340            &stream.next().await.ok_or("second batch missing")??,
1341        )];
1342        while let Some(batch) = stream.next().await {
1343            ids.push(batch_id(&batch?));
1344        }
1345        assert_eq!(ids, [2, 10, 20]);
1346        assert_eq!(metrics.snapshot().scheduler_batches_emitted, 4);
1347        assert_eq!(metrics.snapshot().scan_partitions_completed, 2);
1348        assert_eq!(limiter.active_file_reads(), 0);
1349        Ok(())
1350    }
1351
1352    #[tokio::test]
1353    async fn merged_stream_drop_cancels_blocked_partitions_and_releases_permits()
1354    -> Result<(), Box<dyn std::error::Error>> {
1355        let ControlledMerge {
1356            mut stream,
1357            limiter,
1358            cancellation,
1359            metrics,
1360            ..
1361        } = controlled_merge()?;
1362
1363        let first = stream.next().await.ok_or("first batch missing")??;
1364        assert_eq!(batch_id(&first), 1);
1365        wait_for_batches(&metrics, 2).await;
1366        assert_eq!(limiter.active_file_reads(), 2);
1367        drop(stream);
1368
1369        assert!(cancellation.is_cancelled());
1370        timeout(Duration::from_secs(5), async {
1371            while limiter.active_file_reads() != 0 {
1372                tokio::task::yield_now().await;
1373            }
1374        })
1375        .await?;
1376        assert_eq!(metrics.snapshot().scheduler_batches_emitted, 2);
1377        assert_eq!(metrics.snapshot().scan_partitions_completed, 0);
1378        Ok(())
1379    }
1380
1381    #[tokio::test]
1382    async fn merged_stream_forwards_one_concurrent_error_and_releases_permits()
1383    -> Result<(), Box<dyn std::error::Error>> {
1384        let options = execution_options()?;
1385        let limiter = ScanReadLimiter::new(options, 2, 2);
1386        let cancellation = ScanCancellation::new();
1387        let metrics = metrics();
1388        let executor: FileExecutor<i32, FileBatchStream> = Arc::new(|task, permit, _| {
1389            async move {
1390                Ok(if task == 1 {
1391                    Box::pin(stream::once(async move {
1392                        let _permit = permit;
1393                        pending::<Result<RecordBatch, crate::DeltaReaderError>>().await
1394                    })) as FileBatchStream
1395                } else {
1396                    Box::pin(stream::once(async move {
1397                        let _permit = permit;
1398                        Err(InvalidConfigurationSnafu {
1399                            reason: "controlled_partition_failure",
1400                        }
1401                        .build())
1402                    })) as FileBatchStream
1403                })
1404            }
1405            .boxed()
1406        });
1407        let admission = Arc::new(|_: &i32| Ok(FileAdmissionDecision::Admit));
1408        let first = PartitionStream::new(
1409            vec![1],
1410            limiter.partition(0)?,
1411            options,
1412            admission.clone(),
1413            Arc::clone(&executor),
1414            metrics.clone(),
1415            cancellation.clone(),
1416        );
1417        let second = PartitionStream::new(
1418            vec![2],
1419            limiter.partition(1)?,
1420            options,
1421            admission,
1422            executor,
1423            metrics.clone(),
1424            cancellation.clone(),
1425        );
1426        let mut stream = direct_stream(VecDeque::from([first, second]), metrics);
1427
1428        let error = timeout(Duration::from_secs(5), stream.next())
1429            .await?
1430            .ok_or("error item missing")?
1431            .expect_err("controlled partition must fail");
1432        assert_eq!(error.code(), "invalid_configuration");
1433        assert!(stream.next().await.is_none());
1434        assert!(cancellation.is_cancelled());
1435        timeout(Duration::from_secs(5), async {
1436            while limiter.active_file_reads() != 0 {
1437                tokio::task::yield_now().await;
1438            }
1439        })
1440        .await?;
1441        Ok(())
1442    }
1443}