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