andygrove commented on code in PR #5365:
URL: https://github.com/apache/datafusion-comet/pull/5365#discussion_r4083936869


##########
native/core/src/execution/delta_dv.rs:
##########
@@ -0,0 +1,2192 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Delta Lake deletion-vector decoding and translation into DataFusion
+//! [`ParquetAccessPlan`]s (feature = "delta").
+//!
+//! Wire formats implemented here (from delta-spark's `DeletionVectorStore` /
+//! `RoaringBitmapArray`, v3.3.2):
+//! - On-disk DV file: 1 version byte at the start of the file; at
+//!   `descriptor.offset`: `[i32 BE size][data: size bytes][i32 BE 
CRC32(data)]`.
+//! - `data`: `[i32 LE magic]` then either
+//!   - magic 1681511376 ("native"): `[i32 LE count]`, then per bitmap
+//!     `[i32 LE size][standard 32-bit RoaringBitmap]`, keys implicit (index);
+//!   - magic 1681511377 ("portable", the spec's 64-bit extension): `[i64 LE
+//!     count]`, then per bitmap `[i32 LE key][standard 32-bit RoaringBitmap]`
+//!     with keys ascending -- exactly [`RoaringTreemap`]'s serialized form.
+
+use std::mem::size_of;
+use std::sync::Arc;
+
+use datafusion::datasource::listing::PartitionedFile;
+use 
datafusion::datasource::physical_plan::parquet::metadata::DFParquetMetadata;
+use datafusion::datasource::physical_plan::parquet::{ParquetAccessPlan, 
RowGroupAccess};
+use datafusion::execution::memory_pool::{MemoryConsumer, MemoryReservation};
+use datafusion::execution::runtime_env::RuntimeEnv;
+use futures::{StreamExt, TryStreamExt};
+use object_store::path::Path;
+use object_store::{ObjectStore, ObjectStoreExt};
+use parquet::arrow::arrow_reader::{RowSelection, RowSelector};
+use parquet::file::metadata::{PageIndexPolicy, ParquetMetaData};
+use roaring::{RoaringBitmap, RoaringTreemap};
+
+use crate::execution::operators::ExecutionError;
+use crate::execution::operators::ExecutionError::GeneralError;
+use datafusion_comet_proto::spark_operator::DeltaSparkDvDescriptor;
+
+const NATIVE_MAGIC: i32 = 1681511376;
+const PORTABLE_MAGIC: i32 = 1681511377;
+
+/// Unframe a DV blob read from `descriptor.offset` of a DV file:
+/// `[i32 BE size][data][i32 BE crc]`. Verifies both the size against the
+/// descriptor's `size_in_bytes` and the CRC32 checksum.
+/// An inline payload carries no framing, so its length is checked against the 
descriptor here,
+/// the way `unframe_dv_blob` checks an on-disk blob's size header.
+fn check_inline_payload_size(
+    file_path: &str,
+    payload: &[u8],
+    size_in_bytes: i32,
+) -> Result<(), ExecutionError> {
+    if payload.len() as i64 != i64::from(size_in_bytes) {
+        return Err(GeneralError(format!(
+            "Inline deletion vector for {file_path} has {} bytes but its 
descriptor says {size_in_bytes}",
+            payload.len()
+        )));
+    }
+    Ok(())
+}
+
+pub fn unframe_dv_blob(blob: &[u8], expected_size: usize) -> Result<&[u8], 
ExecutionError> {
+    if blob.len() < 8 {
+        return Err(GeneralError(format!(
+            "Deletion vector blob too short: {} bytes",
+            blob.len()
+        )));
+    }
+    let size = i32::from_be_bytes(blob[0..4].try_into().unwrap());
+    if size < 0 || size as usize != expected_size {
+        return Err(GeneralError(format!(
+            "Deletion vector size mismatch: file says {size}, descriptor says 
{expected_size}"
+        )));
+    }
+    let end = 4 + size as usize;
+    if blob.len() < end + 4 {
+        return Err(GeneralError(format!(
+            "Deletion vector blob truncated: need {} bytes, have {}",
+            end + 4,
+            blob.len()
+        )));
+    }
+    let data = &blob[4..end];
+    let expected_crc = i32::from_be_bytes(blob[end..end + 
4].try_into().unwrap());
+    let actual_crc = crc32fast::hash(data) as i32;
+    if expected_crc != actual_crc {
+        return Err(GeneralError(
+            "Deletion vector checksum mismatch".to_string(),
+        ));
+    }
+    Ok(data)
+}
+
+/// Deserialize the magic-prefixed RoaringBitmapArray into a 64-bit treemap of
+/// deleted row indexes.
+pub fn deserialize_dv_bitmap(data: &[u8]) -> Result<RoaringTreemap, 
ExecutionError> {
+    if data.len() < 4 {
+        return Err(GeneralError(
+            "Deletion vector bitmap too short for magic number".to_string(),
+        ));
+    }
+    let magic = i32::from_le_bytes(data[0..4].try_into().unwrap());
+    let rest = &data[4..];
+    match magic {
+        PORTABLE_MAGIC => RoaringTreemap::deserialize_from(rest)
+            .map_err(|e| GeneralError(format!("Invalid portable deletion 
vector bitmap: {e}"))),
+        NATIVE_MAGIC => {
+            if rest.len() < 4 {
+                return Err(GeneralError(
+                    "Native deletion vector bitmap missing count".to_string(),
+                ));
+            }
+            let count = i32::from_le_bytes(rest[0..4].try_into().unwrap());
+            if count < 0 {
+                return Err(GeneralError(format!(
+                    "Invalid RoaringBitmapArray length ({count} < 0)"
+                )));
+            }
+            let mut pos = 4usize;
+            let mut treemap = RoaringTreemap::new();
+            for key in 0..count as u64 {
+                if rest.len() < pos + 4 {
+                    return Err(GeneralError(
+                        "Native deletion vector bitmap truncated".to_string(),
+                    ));
+                }
+                let size = i32::from_le_bytes(rest[pos..pos + 
4].try_into().unwrap());
+                pos += 4;
+                if size < 0 || rest.len() < pos + size as usize {
+                    return Err(GeneralError(
+                        "Native deletion vector bitmap truncated".to_string(),
+                    ));
+                }
+                let bitmap = RoaringBitmap::deserialize_from(&rest[pos..pos + 
size as usize])
+                    .map_err(|e| {
+                        GeneralError(format!("Invalid deletion vector 
sub-bitmap: {e}"))
+                    })?;
+                pos += size as usize;
+                for value in bitmap {
+                    treemap.insert((key << 32) | value as u64);
+                }
+            }
+            Ok(treemap)
+        }
+        other => Err(GeneralError(format!(
+            "Unexpected RoaringBitmapArray magic number {other}"
+        ))),
+    }
+}
+
+/// Translate deleted row indexes into a [`ParquetAccessPlan`]: fully-deleted
+/// row groups become `Skip`, untouched groups stay `Scan`, and partially
+/// deleted groups get a `RowSelection` selecting the complement of the deleted
+/// rows. Page-index pruning later INTERSECTS with these selections, so DV
+/// skips and page skips compose.
+pub fn build_access_plan(
+    row_group_row_counts: &[i64],
+    deleted: &RoaringTreemap,
+) -> Result<ParquetAccessPlan, ExecutionError> {
+    let mut plan = ParquetAccessPlan::new_all(row_group_row_counts.len());
+    // Single sweep over the (sorted) deleted row indexes, bucketing by row 
group.
+    let mut deleted_iter = deleted.iter().peekable();
+    let mut group_start = 0u64;
+    for (idx, &num_rows) in row_group_row_counts.iter().enumerate() {
+        // A corrupt footer can report a negative row count. `num_rows as u64` 
would otherwise
+        // wrap it into a huge positive value, silently corrupting every 
row-group boundary
+        // computed from `group_start`/`group_end` below (and therefore which 
deleted row indexes
+        // land in which row group) instead of failing loudly.
+        if num_rows < 0 {
+            return Err(GeneralError(format!(
+                "Parquet footer reports a negative row count ({num_rows}) for 
row group {idx}"
+            )));
+        }
+        let num_rows = num_rows as u64;
+        let group_end = group_start.checked_add(num_rows).ok_or_else(|| {
+            GeneralError(format!(
+                "Parquet footer row counts overflow at row group {idx} 
({group_start} + {num_rows})"
+            ))
+        })?;
+        let mut selectors: Vec<RowSelector> = Vec::new();
+        let mut cursor = group_start;
+        let mut deleted_in_group = 0u64;
+        while let Some(&row) = deleted_iter.peek() {
+            if row >= group_end {
+                break;
+            }
+            deleted_iter.next();
+            deleted_in_group += 1;
+            if row > cursor {
+                selectors.push(RowSelector::select((row - cursor) as usize));
+            }
+            // Merge runs of consecutive deleted rows into one skip.
+            match selectors.last_mut() {
+                Some(last) if last.skip => last.row_count += 1,
+                _ => selectors.push(RowSelector::skip(1)),
+            }
+            cursor = row + 1;
+        }
+        if deleted_in_group == num_rows && num_rows > 0 {
+            plan.skip(idx);
+        } else if deleted_in_group > 0 {
+            if group_end > cursor {
+                selectors.push(RowSelector::select((group_end - cursor) as 
usize));
+            }
+            plan.scan_selection(idx, RowSelection::from(selectors));
+        }
+        group_start = group_end;
+    }
+    // A deleted index beyond the file's total row count means the DV does not
+    // belong to this file (stale or corrupted metadata); silently dropping it
+    // would under-apply deletions.
+    if let Some(&row) = deleted_iter.peek() {
+        return Err(GeneralError(format!(
+            "Deletion vector marks row {row} but the file only has 
{group_start} rows"
+        )));
+    }
+    Ok(plan)
+}
+
+/// Verify a decoded deletion vector's row count matches the descriptor's
+/// declared `cardinality`, mirroring Delta's JVM reader
+/// (`StoredBitmap.validateCardinality`). The CRC and framing checks catch
+/// corruption but not a stale, otherwise well-formed bitmap whose row count
+/// no longer matches the descriptor -- that would silently under- or
+/// over-delete rows.
+fn validate_cardinality(
+    file_path: &str,
+    expected: i64,
+    deleted: &RoaringTreemap,
+) -> Result<(), ExecutionError> {
+    let actual = deleted.len();
+    if actual != expected as u64 {
+        return Err(GeneralError(format!(
+            "Deletion vector for {file_path} has cardinality mismatch: 
descriptor says {expected}, decoded bitmap has {actual} deleted rows"
+        )));
+    }
+    Ok(())
+}
+
+/// One data file plus everything needed to apply its deletion vector. The
+/// file's size comes from `file.object_meta.size` (built by the planner from
+/// the proto's `file_size`).
+///
+/// `data_store` and `dv_store` are resolved by the caller *before* entering
+/// the async `attach_access_plans` runtime (see its doc comment): building an
+/// object store is sync I/O that, for a cold S3 authority, internally issues
+/// its own `Handle::block_on` calls, which panics if nested inside another
+/// `block_on`. Resolving up front means this module never constructs a
+/// store itself.
+pub struct DvScanFile {
+    pub file: PartitionedFile,
+    /// Full URL of the data file (proto `file_path`).
+    pub file_path: String,
+    pub dv: Option<DeltaSparkDvDescriptor>,
+    /// Object store for `file_path`, pre-resolved by the caller. Only read
+    /// when `dv` is `Some` (files without a deletion vector never open their
+    /// footer here), but every file carries one so the struct's shape
+    /// doesn't depend on whether a deletion vector is present.
+    pub data_store: Arc<dyn ObjectStore>,
+    /// Store and within-store path for an on-disk deletion vector's absolute
+    /// path, pre-resolved by the caller. `None` when the file has no
+    /// deletion vector or the deletion vector is stored inline.
+    pub dv_store: Option<(Arc<dyn ObjectStore>, Path)>,
+}
+
+/// Execution-memory-pool reservation covering one file's expanded DV row 
selectors across
+/// their *entire* lifetime attached to a scan -- from `build_access_plan`'s 
construction
+/// through DataFusion 54.1's reader normalizing the attached 
[`ParquetAccessPlan`]
+/// (`create_initial_plan`'s deep clone plus `into_overall_row_selection`'s 
combined
+/// `RowSelection`; see [`reader_peak_bytes`]) -- attached to the file's 
[`PartitionedFile`]
+/// extensions alongside its [`ParquetAccessPlan`]. The reservation's lifetime 
is tied to the
+/// `PartitionedFile` it is attached to, so it is released back to the pool 
exactly when the
+/// plan is dropped (query completion or an early-terminated scan), never held 
open longer.
+/// Newtype-wrapped so it occupies its own slot in the multi-slot, type-keyed 
`extensions` map
+/// (`datafusion_common::extensions::Extensions`) alongside the plan, rather 
than a bare
+/// `MemoryReservation` colliding with one some other extension might attach.
+pub struct DvAccessPlanReservation(pub MemoryReservation);
+
+/// Total number of [`RowSelector`]s materialized across `plan`'s per-row-group
+/// selections (`RowGroupAccess::Selection`); `Scan`/`Skip` row groups
+/// contribute none. An alternating deleted/retained bitmap produces one
+/// non-coalescing selector per row (see [`reader_peak_bytes`]'s doc comment
+/// for the worst-case accounting), so this count -- not the deletion
+/// vector's cardinality -- is the thing that must be bounded and reserved
+/// against the execution memory pool.
+fn total_selectors(plan: &ParquetAccessPlan) -> usize {
+    plan.inner()
+        .iter()
+        .map(|access| match access {
+            RowGroupAccess::Selection(selection) => selection.iter().count(),
+            _ => 0,
+        })
+        .sum()
+}
+
+/// Multiplier bounding the peak allocation live *during construction* of one
+/// file's [`RowSelection`]s, relative to the conservative selector-count
+/// bound `S = 2 * cardinality + num_row_groups` (one non-coalescing selector
+/// per deleted row in the worst-case alternating pattern, doubled, plus up to
+/// one extra boundary selector per row group). Split `S` into `r`, the
+/// selectors already retained from row groups `build_access_plan` has
+/// finished, and `c`, the selectors accumulated so far in the current row
+/// group's source `Vec`; `r` and `c` partition the selectors counted toward
+/// `S`, so `r + c <= S` always. While the current group is being built, the
+/// `Vec`'s doubling growth strategy can leave its backing allocation at up to
+/// `2 * c` (the next power-of-two capacity above `c`). Once the group
+/// finishes, `RowSelection::from(Vec)` (parquet's `FromIterator` impl,
+/// `with_capacity` + copy) builds a second, separate `Vec` of size `c` from
+/// that source while the source is still alive, so at the moment the copy
+/// begins, the retained selectors, the current group's doubled source `Vec`,
+/// and the copy are all live simultaneously: `r + 2c + c = r + 3c`. Since
+/// `r >= 0`, `r + 3c <= 3r + 3c = 3(r + c) <= 3S`. 3x covers that peak.
+const CONSTRUCTION_PEAK_FACTOR: usize = 3;
+
+/// Upper bound on how much larger a `Vec`'s backing allocation can be than 
its element count
+/// after being built by repeated pushes: `std`'s doubling growth strategy 
never leaves a `Vec`
+/// of `n` elements with a backing allocation larger than the next power of 
two above `n`, which
+/// is at most `2 * n` for any `n >= 1`.
+const VEC_GROWTH_CAPACITY_FACTOR: usize = 2;
+
+/// `RawVec`'s minimum non-zero capacity for element sizes `<= 1024` bytes 
([`RowSelector`] is
+/// 16 bytes on 64-bit platforms: a `usize` row count plus a padded `bool`). 
Applied once per
+/// row group (or per contiguous run of row groups) a fresh 
`from_fn`/`FlatMap`-driven `Vec`
+/// gets built for (see [`reader_peak_bytes`]), so even a group or run whose 
true selector count
+/// is tiny still pays this floor.
+const MIN_VEC_CAPACITY_SELECTORS: usize = 4;
+
+/// Conservative upper bound, in bytes, on the peak allocation live while 
DataFusion 54.1's
+/// reader normalizes one file's attached [`ParquetAccessPlan`] -- the 
allocation this module's
+/// steady-state reservation must cover, not merely the plan's own retained 
selector bytes.
+/// THREE allocations can be live simultaneously by the time 
`into_overall_row_selection`
+/// returns, not two -- the clone is only exact when page-index pruning never 
touches it:
+///
+/// 1. **Attached original** (`selectors`, exact): `create_initial_plan` 
deep-clones the
+///    attached plan while the original remains reachable from the file's 
`extensions` until
+///    the scan consumes it. The ORIGINAL's own selector `Vec`s are exact -- a 
coalesced
+///    [`RowSelection`] built via `RowSelection::from(Vec<RowSelector>)` (what
+///    `build_access_plan` uses) has no excess capacity, because that 
conversion is a plain
+///    `with_capacity(len)` copy, not a `size_hint`-blind fold.
+/// 2. **The clone, possibly capacity-inflated** (`<= 
VEC_GROWTH_CAPACITY_FACTOR * selectors +
+///    MIN_VEC_CAPACITY_SELECTORS * num_row_groups`): if page-index pruning 
fires
+///    (`PagePruningAccessPlanFilter`; `access_plan.rs`'s `scan_selection` on 
a row group that
+///    already carries a `RowGroupAccess::Selection` calls 
`existing.intersection(&page_derived)`
+///    -- `RowSelection::intersection` -> `intersect_row_selections`), it 
replaces the CLONE's
+///    per-row-group selection with that intersection's output. 
`intersect_row_selections` is
+///    ANOTHER `from_fn` generator with `size_hint() == (0, None)`, so each 
intersected row
+///    group's backing `Vec` starts at `with_capacity(0)` and doubles as it 
grows, independent
+///    of whatever capacity the pre-intersection selection had. This inflated 
clone is still
+///    live when `into_overall_row_selection` later moves its buffer. Term 1's 
exactness
+///    guarantee holds for the ORIGINAL always, and for the clone only when 
page-index pruning
+///    never fires against it -- once it does, the clone must be charged at 
the SAME
+///    growth-capped bound as a fresh combined-selection `Vec` (term 3), 
summed once per row
+///    group rather than once per run, since each row group's `Selection` is 
intersected
+///    independently.
+/// 3. **Per-run combined-selection allocation** (`<= 
VEC_GROWTH_CAPACITY_FACTOR * (selectors +
+///    num_row_groups) + MIN_VEC_CAPACITY_SELECTORS * num_row_groups`): 
`into_overall_row_selection`
+///    collects each contiguous run of row groups' selectors into a *new* 
`RowSelection` via a
+///    `FlatMap` whose `size_hint().0 == 0`, so that run's `Vec` starts at 
`with_capacity(0)`
+///    and doubles as it grows -- capping its backing allocation at
+///    `max(MIN_VEC_CAPACITY_SELECTORS, next_power_of_two(len))`, which is at 
most
+///    `MIN_VEC_CAPACITY_SELECTORS + VEC_GROWTH_CAPACITY_FACTOR * len` for a 
run of `len`
+///    selectors. `len` is at most that run's share of `selectors` plus one 
boundary selector
+///    per `RowGroupAccess::Scan` row group in the run (`Scan` always 
contributes exactly one
+///    `RowSelector::select(num_rows)`; see `access_plan.rs`'s 
`into_overall_row_selection`).
+///    Summing across at most `num_row_groups` runs (each spans >= 1 row 
group) bounds the total
+///    at `VEC_GROWTH_CAPACITY_FACTOR * selectors + 
(MIN_VEC_CAPACITY_SELECTORS +
+///    VEC_GROWTH_CAPACITY_FACTOR) * num_row_groups`.
+///
+/// Summing all three terms and converting to bytes: `((1 + 2 * 
VEC_GROWTH_CAPACITY_FACTOR) *
+/// selectors + (2 * MIN_VEC_CAPACITY_SELECTORS + VEC_GROWTH_CAPACITY_FACTOR) 
* num_row_groups)
+/// * size_of::<RowSelector>()` -- with the constants above, `(5 * selectors + 
10 *
+///   num_row_groups) * size_of::<RowSelector>()`. Checked against two 
measured worst cases:
+///
+/// - No page-index pruning (the original P2 report; term 2 stays exact): one 
2,000,000-row
+///   group, 1,000,000 alternating deletions, `selectors = 2,000,000`. 
Measured allocator peak
+///   97,554,457 B; the byte-for-byte accounting for the attached original 
plus the (here,
+///   exact) clone plus the inflated combined selection explains 97,554,432 B 
of that, a 25 B
+///   residue we did not attribute. This bound gives 160,000,160 B -- much 
looser here because
+///   it must also cover the next case, where the clone is NOT exact.
+/// - Page-index pruning fires against the clone: one 1,048,577-row group, 
`selectors =
+///   1,048,577`. Measured peak 83,886,096 B; this bound gives 83,886,320 B (a 
224 B, <1%
+///   margin -- deliberately tight, since this is the case that drives the 
bound).
+///
+/// Uses checked arithmetic throughout: a selector or row-group count large 
enough to overflow
+/// `usize` indicates a corrupted or malicious input, reported as a clean 
error rather than
+/// panicking.
+fn reader_peak_bytes(selectors: usize, num_row_groups: usize) -> Result<usize, 
ExecutionError> {
+    let overflow = || {
+        GeneralError(format!(
+            "Deletion vector reader-peak bound overflowed for {selectors} 
selectors and \
+             {num_row_groups} row groups"
+        ))
+    };
+    // Term 1: the attached original -- exact, untouched by page-index pruning 
(only the clone
+    // is ever intersected; see the doc comment above).
+    let attached_term = selectors;
+    // Term 2: the clone, bounded as if page-index pruning DID fire against 
every row group
+    // (safe even when it doesn't: term 2's bound is always >= `selectors`, so 
it never
+    // undershoots the exact case either).
+    let clone_growth = selectors
+        .checked_mul(VEC_GROWTH_CAPACITY_FACTOR)
+        .ok_or_else(overflow)?;
+    let clone_floor = num_row_groups
+        .checked_mul(MIN_VEC_CAPACITY_SELECTORS)
+        .ok_or_else(overflow)?;
+    let clone_term = 
clone_growth.checked_add(clone_floor).ok_or_else(overflow)?;
+    // Term 3: into_overall_row_selection's per-run combined-selection 
allocation.
+    let combined_growth = selectors
+        .checked_mul(VEC_GROWTH_CAPACITY_FACTOR)
+        .ok_or_else(overflow)?;
+    let combined_floor = num_row_groups
+        .checked_mul(MIN_VEC_CAPACITY_SELECTORS + VEC_GROWTH_CAPACITY_FACTOR)
+        .ok_or_else(overflow)?;
+    let combined_term = combined_growth
+        .checked_add(combined_floor)
+        .ok_or_else(overflow)?;
+
+    let selector_bound = attached_term
+        .checked_add(clone_term)
+        .and_then(|sum| sum.checked_add(combined_term))
+        .ok_or_else(overflow)?;
+    selector_bound
+        .checked_mul(size_of::<RowSelector>())
+        .ok_or_else(overflow)
+}
+
+/// Upper bound, in [`RowSelector`]s, on how many extra selectors the parquet 
reader's
+/// page-index pruning can add on top of the deletion vector's own selection 
when normalizing
+/// one file, from that file's already-fetched [`ParquetMetaData`].
+///
+/// `intersect_row_selections` (parquet's `selection.rs`), which combines a 
page-pruning
+/// selection with the deletion vector's selection, is a `from_fn` generator 
whose
+/// `size_hint()` is `(0, None)`: for inputs of length `a` and `b`, its output 
can have up to
+/// `a + b` selectors -- longer than either input. Bounding the page-pruning 
side of that sum
+/// requires knowing how many selectors a page-index-derived selection could 
produce: at most
+/// two per data page (one skip, one select, in the worst case of alternating 
page-level
+/// pruning decisions), summed over every column of every row group.
+///
+/// Returns `0` when `metadata` carries no offset index 
(`metadata.offset_index()` is `None`).
+/// This is provably safe, not merely a convenient default: page-index pruning 
cannot produce a
+/// page-level selection without the offset index to locate pages by, so there 
are no
+/// page-pruning selectors to bound. The offset index is fetched with
+/// `PageIndexPolicy::Optional` from the same `FileMetadataCache` entry the 
scan's reader later
+/// reopens (see [`attach_access_plan`]'s footer-fetch comment), so this 
function observes
+/// exactly what the reader will see.
+///
+/// Uses checked arithmetic throughout for the same reason as 
[`admission_bound_bytes`].
+fn page_selection_bound_selectors(metadata: &ParquetMetaData) -> Result<usize, 
ExecutionError> {
+    let Some(offset_index) = metadata.offset_index() else {
+        return Ok(0);
+    };
+    let overflow = || {
+        GeneralError(
+            "Deletion vector page-selection bound overflowed while summing 
offset-index page \
+             locations"
+                .to_string(),
+        )
+    };
+    let mut total_page_locations = 0usize;
+    for row_group in offset_index {
+        for column in row_group {
+            total_page_locations = total_page_locations
+                .checked_add(column.page_locations().len())
+                .ok_or_else(overflow)?;
+        }
+    }
+    total_page_locations.checked_mul(2).ok_or_else(overflow)
+}
+
+/// Execution-memory-pool admission bound, in bytes, for one file's 
deletion-vector access
+/// plan -- reserved *before* calling `build_access_plan` (see 
[`attach_access_plan`]'s
+/// pre-reserve call site) to cover the larger of two peaks live at different 
points in the
+/// plan's lifetime. In practice the reader-normalization peak below dominates 
the construction
+/// peak unconditionally for any non-trivial input (`reader_peak_bytes(S, G) = 
(5S + 10G) *
+/// size_of::<RowSelector>()` always exceeds `CONSTRUCTION_PEAK_FACTOR * S *
+/// size_of::<RowSelector>() = 3S * size_of::<RowSelector>()` once `S >= 1`, 
since the `5S` term
+/// alone already exceeds `3S`); the construction term is retained as a 
documented floor rather
+/// than dropped, since it is cheap to compute and keeps this bound correct 
even if the reader's
+/// growth factors ever shrink below construction's.
+///
+/// - **Construction peak** (`CONSTRUCTION_PEAK_FACTOR * S`, see that 
constant's doc comment):
+///   live while `build_access_plan` builds the plan's `RowSelection`s. 
Construction's
+///   transient allocations fully unwind before `build_access_plan` returns, 
so this peak never
+///   overlaps the reader-normalization peak below.
+/// - **Reader-normalization peak** (`reader_peak_bytes(S + 
page_bound_selectors,
+///   num_row_groups)`, see that function): live later, once DataFusion's 
reader normalizes the
+///   attached plan. `S = 2 * cardinality + num_row_groups` is the same 
conservative bound on
+///   the plan's final retained selector count used for the construction peak 
-- it provably
+///   bounds `R = total_selectors(&plan)` (`R <= S`, from `build_access_plan`'s
+///   one-non-coalescing-selector-per-deleted-row worst case plus one boundary 
selector per row
+///   group), so `S + page_bound_selectors` bounds `R` after page-index 
inflation the same way
+///   `S` bounds `R` before it.
+///
+/// These two peaks never overlap in time, so `max` -- not `sum` -- is the 
correct combinator:
+/// reserving their sum would over-reserve for no safety benefit.
+///
+/// Deliberately not clamped by the file's total row count here, unlike the 
reader-peak target
+/// `attach_access_plan` resizes down to after construction (see that call 
site): `S`'s
+/// `+ num_row_groups` boundary term is a worst-case padding margin that can 
legitimately exceed
+/// the total row count for a small, heavily-deleted file, and admission 
sizing has no actual
+/// retained-selector count yet to clamp against -- only after construction, 
once `R` is known,
+/// is clamping to the total row count both meaningful and strictly tighter. 
Leaving this bound
+/// unclamped only ever makes admission more conservative, never less safe.
+///
+/// Uses checked arithmetic throughout: a cardinality, row-group count, or 
page bound large
+/// enough to overflow `usize` while computing this bound indicates a 
corrupted or malicious
+/// descriptor, reported as a clean error rather than panicking.
+fn admission_bound_bytes(
+    cardinality: i64,
+    num_row_groups: usize,
+    page_bound_selectors: usize,
+) -> Result<usize, ExecutionError> {
+    let overflow = || {
+        GeneralError(format!(
+            "Deletion vector admission bound overflowed for cardinality 
{cardinality}, \
+             {num_row_groups} row groups, and page bound 
{page_bound_selectors} selectors"
+        ))
+    };
+    let cardinality_usize = usize::try_from(cardinality).map_err(|_| 
overflow())?;
+    // S: the conservative bound on the plan's final *retained* selector count 
(what
+    // `total_selectors(&plan)` cannot exceed) -- unchanged from the 
pre-existing
+    // construction-only bound this function replaces.
+    let s = cardinality_usize
+        .checked_mul(2)
+        .and_then(|doubled| doubled.checked_add(num_row_groups))
+        .ok_or_else(overflow)?;
+
+    let construction_bytes = s
+        .checked_mul(size_of::<RowSelector>())
+        .and_then(|bytes| bytes.checked_mul(CONSTRUCTION_PEAK_FACTOR))
+        .ok_or_else(overflow)?;
+
+    let s_plus_page = 
s.checked_add(page_bound_selectors).ok_or_else(overflow)?;
+    let reader_bytes = reader_peak_bytes(s_plus_page, num_row_groups)?;
+
+    Ok(construction_bytes.max(reader_bytes))
+}
+
+/// Upper bound on concurrent DV-blob and footer fetches per partition. Both
+/// are small ranged reads, so a modest fan-out hides object-store latency
+/// without flooding the store client.
+const DV_FETCH_CONCURRENCY: usize = 8;
+
+/// Called via `block_on` at plan-creation time on the executor task: DV blobs
+/// are small ranged reads and footers are needed to learn row-group
+/// boundaries. Files are fetched concurrently (bounded by
+/// [`DV_FETCH_CONCURRENCY`]) with input order preserved. Footer fetches go
+/// through the scan's shared FileMetadataCache, so the scan's subsequent open
+/// of the same file is served from cache. That reuse relies on each input
+/// [`PartitionedFile`] being returned as-is (only `with_extension` applied),
+/// never rebuilt: the cache entry is keyed by this exact `object_meta` and the
+/// scan later looks it up through the same struct.
+///
+/// Deliberately takes no object-store options map and imports no
+/// store-construction helper: every [`DvScanFile`] arrives with its stores
+/// already resolved by the caller (see its doc comment), so this async path
+/// structurally cannot build an object store -- only `runtime_env` is still
+/// threaded through, for the shared `FileMetadataCache` and (per file) the
+/// execution `MemoryPool` each expanded access plan's row selectors are
+/// reserved against -- see [`DvAccessPlanReservation`].
+pub async fn attach_access_plans(
+    runtime_env: Arc<RuntimeEnv>,
+    files: Vec<DvScanFile>,
+) -> Result<Vec<PartitionedFile>, ExecutionError> {
+    futures::stream::iter(files)
+        .map(|scan_file| attach_access_plan(Arc::clone(&runtime_env), 
scan_file))
+        .buffered(DV_FETCH_CONCURRENCY)
+        .try_collect()
+        .await
+}
+
+/// Resolve one file's deletion vector into an attached [`ParquetAccessPlan`];
+/// files without a DV pass through untouched.
+async fn attach_access_plan(
+    runtime_env: Arc<RuntimeEnv>,
+    scan_file: DvScanFile,
+) -> Result<PartitionedFile, ExecutionError> {
+    let DvScanFile {
+        file,
+        file_path,
+        dv,
+        data_store,
+        dv_store,
+    } = scan_file;
+    let dv = match dv {
+        Some(dv) => dv,
+        None => return Ok(file),
+    };
+    // Delta's canonical `DeletionVectorDescriptor.EMPTY`: inline storage, 
empty
+    // payload, size 0, cardinality 0. Spark's reader returns all rows for it;
+    // decoding would fail (the empty payload is too short for a magic
+    // number), so pass the file through unchanged before attempting to read 
it.
+    if dv.cardinality == 0 && dv.size_in_bytes == 0 {
+        return Ok(file);
+    }
+    if dv.size_in_bytes < 0 {
+        return Err(GeneralError(format!(
+            "Deletion vector for {file_path} has negative size {}",
+            dv.size_in_bytes
+        )));
+    }
+    if dv.cardinality < 0 {
+        return Err(GeneralError(format!(
+            "Deletion vector for {file_path} has negative cardinality {}",
+            dv.cardinality
+        )));
+    }
+
+    let data: Vec<u8> = if let Some(inline) = dv.inline_data {
+        check_inline_payload_size(&file_path, &inline, dv.size_in_bytes)?;
+        inline
+    } else if let Some(dv_path) = &dv.absolute_path {
+        let offset = dv
+            .offset
+            .ok_or_else(|| GeneralError("On-disk deletion vector missing 
offset".into()))?;
+        if offset < 0 {
+            return Err(GeneralError(format!(
+                "Deletion vector for {file_path} has negative offset {offset}"
+            )));
+        }
+        let offset = offset as u64;
+        // [i32 BE size][data: size_in_bytes][i32 BE crc]
+        let framed_len = 4 + dv.size_in_bytes as u64 + 4;
+        let (store, dv_store_path) = dv_store.ok_or_else(|| {
+            GeneralError(format!(
+                "Deletion vector for {file_path} has an absolute path but no 
pre-resolved object store"
+            ))
+        })?;
+        let blob = store
+            .get_range(&dv_store_path, offset..offset + framed_len)
+            .await
+            .map_err(|e| GeneralError(format!("Failed to read deletion vector 
{dv_path}: {e}")))?;
+        unframe_dv_blob(&blob, dv.size_in_bytes as usize)?.to_vec()
+    } else {
+        return Err(GeneralError(
+            "Deletion vector descriptor has neither inline data nor a 
path".into(),
+        ));
+    };
+    let deleted = deserialize_dv_bitmap(&data)
+        .map_err(|e| GeneralError(format!("Invalid deletion vector for 
{file_path}: {e}")))?;
+    validate_cardinality(&file_path, dv.cardinality, &deleted)?;
+
+    // Row-group boundaries come from the data file's footer, fetched through 
the scan's
+    // shared FileMetadataCache with the page index loaded eagerly and the 
scan's metadata
+    // size hint (mirroring EagerPageIndexReaderFactory): the one fetch here 
also serves the
+    // subsequent data-file open, so DV files pay no extra footer round-trip. 
Keyed by
+    // `file.object_meta`, the exact ObjectMeta the scan's reader factory will 
look up.
+    let metadata_cache = runtime_env.cache_manager.get_file_metadata_cache();
+    let metadata = DFParquetMetadata::new(data_store.as_ref(), 
&file.object_meta)
+        .with_file_metadata_cache(Some(metadata_cache))
+        .with_page_index_policy(Some(PageIndexPolicy::Optional))
+        
.with_metadata_size_hint(Some(crate::parquet::parquet_exec::METADATA_SIZE_HINT))
+        .fetch_metadata()
+        .await
+        .map_err(|e| GeneralError(format!("Failed to read parquet footer of 
{file_path}: {e}")))?;
+    let row_counts: Vec<i64> = metadata
+        .row_groups()
+        .iter()
+        .map(|rg| rg.num_rows())
+        .collect();
+
+    // Pre-reserve the admission bound *before* calling build_access_plan: 
this bound covers
+    // both construction's own transient peak AND the larger peak DataFusion's 
reader hits
+    // later while normalizing the attached plan (`create_initial_plan`'s deep 
clone plus
+    // `into_overall_row_selection`'s combined RowSelection) -- see 
admission_bound_bytes and
+    // reader_peak_bytes. Reserving first means a rejection happens before any 
large `Vec` is
+    // allocated, not after -- see reader_peak_bytes's doc comment for the 
measured worst
+    // cases. The error message names this as a construction-phase rejection 
(contains
+    // "construct"), textually distinct from the steady-state message below, 
so callers/logs
+    // can tell which phase failed.
+    let page_bound_selectors = page_selection_bound_selectors(&metadata)?;
+    let admission_bytes =
+        admission_bound_bytes(dv.cardinality, row_counts.len(), 
page_bound_selectors)?;
+    let reservation =
+        
MemoryConsumer::new("DeltaDeletionVectorAccessPlan").register(&runtime_env.memory_pool);

Review Comment:
   Every DV'd file registers its own consumer here. Because the reservation 
lives in the `PartitionedFile` extensions, the registration stays until the 
plan is dropped at the end of the task. The default `fair_unified` pool rejects 
a request once the task's total would exceed `pool_size / num_consumers`, and 
it counts every consumer in the task. So a partition with many DV'd files 
lowers the limit for every other native operator in that task, including a hash 
join build that can't spill.
   
   I reproduced this on Spark 3.5 with `spark.memory.offHeap.size=256m`. I used 
16 files packed into one partition, each with a DV deleting a single row, 
joined to a broadcast of `range(1500000)`. With the native Delta scan it fails 
with `Failed to acquire 24000000 bytes where 11520 bytes already reserved and 
the fair limit is 14913080 bytes, 18 registered`. It passes with the contrib 
off, and it passes natively on the same rows written without DVs. The DV 
reservations add up to 11.5 KB, so the consumer count alone is failing the join.
   
   Could `attach_access_plans` register one consumer per call and hand each 
file `new_empty()` from it? That shares one registration across the partition 
and keeps the per-file grow, resize and release. The existing reservation tests 
all use `GreedyMemoryPool`, which has no consumer-count term. A test pool that 
mirrors `CometFairMemoryPool`'s check, or one that just counts `register` 
calls, would pin this. Related, the `maxDeletedRowsPerFile` doc says the 
selectors are retained for the file's scan. They're actually held until the 
task finishes, for every file in the partition. Could the doc say that the cap 
bounds one file rather than what a task holds?



##########
contrib/delta-spark/src/main/scala/org/apache/spark/sql/comet/CometDeltaNativeScanExec.scala:
##########
@@ -0,0 +1,310 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.spark.sql.comet
+
+import org.apache.spark.rdd.RDD
+import org.apache.spark.sql.catalyst.expressions._
+import org.apache.spark.sql.catalyst.plans.QueryPlan
+import org.apache.spark.sql.catalyst.plans.physical.{Partitioning, 
UnknownPartitioning}
+import org.apache.spark.sql.execution.{FileSourceScanExec, InSubqueryExec, 
ReusedSubqueryExec, ScalarSubquery, SparkPlan, SubqueryAdaptiveBroadcastExec}
+import org.apache.spark.sql.execution.datasources.HadoopFsRelation
+import org.apache.spark.sql.execution.metric.SQLMetric
+import org.apache.spark.sql.types.StructType
+import org.apache.spark.sql.vectorized.ColumnarBatch
+
+import org.apache.comet.contrib.delta.DeltaSparkScanEnvelope
+import org.apache.comet.serde.OperatorOuterClass
+import org.apache.comet.serde.OperatorOuterClass.Operator
+
+/**
+ * Native scan node for Delta Lake tables (contrib). Delta's own planning (log 
replay, snapshot
+ * resolution, partition pruning) has already run inside delta-spark by the 
time this node is
+ * created from the DSv1 [[FileSourceScanExec]]; file listing and split 
planning are delegated to
+ * a [[CometScanExec]] helper, and data reads execute through Comet's native 
DataFusion parquet
+ * machinery, inheriting row-group and page-index pruning.
+ *
+ * DPP: `runtimeFilters` is a constructor field included in equality, so its 
rewrite (via
+ * [[CometScanWithPlanData]]) survives plan copies -- a transient field would 
be dropped by
+ * `TreeNode.makeCopy` on MERGE re-planning (the CometIcebergNativeScanExec 
lesson).
+ */
+case class CometDeltaNativeScanExec(
+    override val nativeOp: Operator,
+    override val output: Seq[Attribute],
+    requiredSchema: StructType,
+    runtimeFilters: Seq[Expression],
+    dataFilters: Seq[Expression],
+    @transient relation: HadoopFsRelation,
+    originalPlan: FileSourceScanExec,
+    override val serializedPlanOpt: SerializedPlan,
+    sourceKey: String)
+    extends CometLeafExec
+    with CometScanWithPlanData {
+
+  override val nodeName: String = s"CometDeltaNativeScan $relation"
+
+  // Derived from (originalPlan, runtimeFilters), never stored: any copy of 
this node
+  // automatically gets a helper consistent with ITS runtimeFilters, avoiding 
the #3510 class of
+  // bug where a stored helper field desyncs from rewritten filters. Costs one 
extra file listing
+  // per executed instance; correctness over the duplicate driver-side listing.
+  //
+  // Forcing invariant: this lazy val is forced by the `metrics` override 
below, and AQE's UI
+  // plan-walk calls `.metrics` on every node MID-PLANNING, including while a 
DPP subquery is
+  // still an adaptive placeholder or a partition filter holds an unresolved 
ScalarSubquery (see
+  // `hasUnevaluableSubqueryFilter` below). That's safe ONLY because 
constructing `scanHelper` is a
+  // cheap case-class build with no file listing, and core's 
`CometScanExec.metrics` touches only
+  // `wrapped.driverMetrics` (populated by Spark's own planning) plus a static 
metric-node
+  // constructor -- neither file listing nor subquery resolution. If core's 
`metrics` ever touches
+  // either, forcing `scanHelper` here would resurrect the AQE mid-planning 
crashes this invariant
+  // prevents.
+  @transient private lazy val scanHelper: CometScanExec =
+    CometDeltaNativeScanExec.planningHelper(originalPlan, runtimeFilters)
+
+  // NOT lazy val: while a DPP subquery is still an adaptive placeholder, or a 
partition filter
+  // holds an unresolved scalar subquery, this returns a temporary value that 
must not be
+  // memoized -- after CometPlanAdaptiveDynamicPruningFilters rewrites the 
filters (DPP case) or
+  // AQE resolves the subquery (scalar case), later reads must see the real 
post-pruning
+  // partition count.
+  override def outputPartitioning: Partitioning =
+    if (hasUnevaluableSubqueryFilter) UnknownPartitioning(0)
+    else UnknownPartitioning(perPartitionData.length)
+
+  // runtimeFilters IS scanHelper.partitionFilters element-for-element, so 
checking runtimeFilters
+  // here avoids constructing/forcing the derived scanHelper just to read 
partitioning. The
+  // InSubqueryExec placeholder shapes mirror
+  // CometPlanAdaptiveDynamicPruningFilters.extractSABData + hasWrappedSAB -- 
keep in sync. The
+  // ScalarSubquery case is probed rather than treated as permanently 
unevaluable: Spark exposes no
+  // public finished/updated flag on ExecSubqueryExpression, but `eval()` 
doubles as one -- it only
+  // reads the cached `result` behind a `require(updated, ...)` guard, while 
the subquery is
+  // actually run by `updateResult()` (invoked separately during prepare/AQE), 
never by `eval()`.
+  // Once resolved, outputPartitioning below reports the real 
perPartitionData.length instead of
+  // staying at zero -- a fused native parent's buildNativeContext requires 
that count to match.
+  private def hasUnevaluableSubqueryFilter: Boolean =
+    runtimeFilters.exists(_.exists {
+      // Match `e: InSubqueryExec` and dispatch on e.plan rather than 
unapplying InSubqueryExec
+      // directly: its unapply arity differs across Spark versions and this 
module ships no
+      // version shim.
+      case e: InSubqueryExec => isAdaptivePlaceholder(e.plan)
+      case s: ScalarSubquery => !isScalarSubqueryResolved(s)
+      case _ => false
+    })
+
+  // `eval()` never triggers the subquery's execution: on a resolved subquery 
it is a pure cached
+  // read of `result` (verified against bytecode: `Predef.require(updated(), 
...)` then a plain
+  // field read), so this probe is safe to call repeatedly, including from 
AQE's mid-planning plan
+  // walks. Pre-resolution, the ONLY throw is `require`'s 
`IllegalArgumentException`; catch exactly
+  // that, since anything else escaping is a genuine bug we must not mask as 
unpartitioned.
+  private def isScalarSubqueryResolved(s: ScalarSubquery): Boolean =
+    try {
+      s.eval()
+      true
+    } catch {
+      case _: IllegalArgumentException => false
+    }
+
+  private def isAdaptivePlaceholder(p: SparkPlan): Boolean = p match {
+    case ReusedSubqueryExec(inner) => isAdaptivePlaceholder(inner)
+    case _: CometSubqueryAdaptiveBroadcastExec => true
+    case _: SubqueryAdaptiveBroadcastExec => true
+    case _ => false
+  }
+
+  override lazy val outputOrdering: Seq[SortOrder] = 
originalPlan.outputOrdering
+
+  override def dynamicPruningFilters: Seq[Expression] = runtimeFilters
+
+  override def withDynamicPruningFilters(filters: Seq[Expression]): SparkPlan 
= {
+    // A real copy: runtimeFilters is a constructor field included in 
equality, so the copy
+    // survives enclosing-block rebuilds, and the derived scanHelper picks up 
the rewritten
+    // filters automatically.
+    copy(runtimeFilters = filters)
+  }
+
+  /**
+   * Lazy split-mode serialization, mirroring CometNativeScanExec: common data 
was serialized at
+   * planning; per-partition file lists serialize here, at execution time.
+   */
+  @transient private lazy val serializedPartitionData
+      : (Array[Byte], Array[Array[Byte]], Array[Seq[String]]) = {
+    // Resolve the helper's DPP subqueries: it holds its own InSubqueryExec 
instances that
+    // Spark's expressions walk does not see (the helper is derived, not a 
child).
+    scanHelper.partitionFilters.foreach {
+      case DynamicPruningExpression(e: InSubqueryExec) if e.values().isEmpty =>
+        e.updateResult()
+      case _ =>
+    }
+
+    val commonBytes = {
+      val deltaScan = DeltaSparkScanEnvelope.unpack(nativeOp)
+      // Scalar subqueries in dataFilters were unresolved at planning; resolve 
them now and
+      // append them as pushed filters, as 
CometNativeScanExec.serializedPartitionData does.
+      // has_data_filters follows their presence, not the serialized count: a 
filter that fails
+      // to serialize still keeps native on the safe timestamp conversion for 
a filtered scan.
+      val resolved = org.apache.comet.contrib.delta.CometDeltaNativeScan
+        .resolvedSubqueryFilters(dataFilters, output, requiredSchema, conf)
+      val common = if (!resolved.hasResolvedFilters) {
+        deltaScan.getCommon
+      } else {
+        val builder = deltaScan.getCommon.toBuilder
+        builder.setHasDataFilters(true)
+        resolved.protos.foreach(builder.addDataFilters)
+        builder.build()
+      }
+      OperatorOuterClass.DeltaSparkScan
+        .newBuilder()
+        .setCommon(common)
+        .setDeltaCommon(deltaScan.getDeltaCommon)
+        .build()
+        .toByteArray
+    }
+
+    val filePartitions = scanHelper.getFilePartitions()
+
+    val tableRoot = 
DeltaSparkScanEnvelope.unpack(nativeOp).getDeltaCommon.getTableRoot
+    val perPartitionBytes = filePartitions.map { filePartition =>
+      org.apache.comet.contrib.delta.CometDeltaNativeScan
+        .serializePartition(filePartition, originalPlan, tableRoot)
+    }.toArray
+
+    val perPartitionPaths = 
filePartitions.map(_.files.map(_.filePath.toString).toSeq).toArray
+
+    (commonBytes, perPartitionBytes, perPartitionPaths)
+  }
+
+  override def commonData: Array[Byte] = serializedPartitionData._1
+
+  override def perPartitionData: Array[Array[Byte]] = 
serializedPartitionData._2
+
+  def perPartitionFilePaths: Array[Seq[String]] = serializedPartitionData._3
+
+  override def doExecuteColumnar(): RDD[ColumnarBatch] = {
+    val nativeMetrics = CometMetricNode.fromCometPlan(this)
+    val serializedPlan = CometExec.serializeNativePlan(nativeOp)
+
+    new CometExecRDD(
+      sparkContext,
+      Seq.empty,
+      Map(sourceKey -> commonData),
+      Map(sourceKey -> perPartitionData),
+      serializedPlan,
+      PlanDataInjector.planFingerprint(serializedPlan),
+      perPartitionData.length,
+      output.length,
+      nativeMetrics,
+      Seq.empty,
+      None,
+      Seq.empty,
+      perPartitionFilePaths = perPartitionFilePaths,
+      reportScanInputMetrics = true)
+  }
+
+  override def doCanonicalize(): CometDeltaNativeScanExec = {
+    val canonOriginal = if (originalPlan != null) {
+      val stripped = originalPlan.copy(partitionFilters =
+        
CometScanUtils.filterUnusedDynamicPruningExpressions(originalPlan.partitionFilters))
+      stripped.doCanonicalize()
+    } else {
+      null
+    }
+    CometDeltaNativeScanExec(
+      nativeOp,
+      output.map(QueryPlan.normalizeExpressions(_, output)),
+      requiredSchema,
+      QueryPlan.normalizePredicates(
+        CometScanUtils.filterUnusedDynamicPruningExpressions(runtimeFilters),
+        output),
+      QueryPlan.normalizePredicates(dataFilters, output),
+      relation,
+      canonOriginal,
+      SerializedPlan(None),
+      "")
+  }
+
+  override def stringArgs: Iterator[Any] = Iterator(output, runtimeFilters)
+
+  override def equals(obj: Any): Boolean = obj match {
+    case other: CometDeltaNativeScanExec =>
+      this.originalPlan == other.originalPlan &&
+      this.serializedPlanOpt == other.serializedPlanOpt &&
+      this.runtimeFilters == other.runtimeFilters &&
+      this.dataFilters == other.dataFilters
+    case _ => false
+  }
+
+  override def hashCode(): Int =
+    java.util.Objects.hash(originalPlan, serializedPlanOpt, runtimeFilters, 
dataFilters)
+
+  private val driverMetricKeys =
+    Set(
+      "numFiles",
+      "filesSize",
+      "numPartitions",
+      "metadataTime",
+      "staticFilesNum",
+      "staticFilesSize",
+      "pruningTime")
+
+  // Forces `scanHelper` (see its doc above for why that -- and reading 
`.metrics` off it -- is
+  // safe even when AQE calls `.metrics` mid-planning against an unresolved 
DPP/scalar subquery).
+  override lazy val metrics: Map[String, SQLMetric] = {
+    CometMetricNode.nativeScanMetrics(session.sparkContext) ++
+      scanHelper.metrics.filter { case (k, _) => driverMetricKeys.contains(k) }
+  }
+}
+
+object CometDeltaNativeScanExec {
+
+  /** File-planning helper: reuses CometScanExec's listing/splitting/DPP 
machinery. */
+  def planningHelper(

Review Comment:
   Claimed DV scans do get split today. The claim requires 
`optimizationsEnabled`, `DeltaParquetFileFormat.isSplitable` returns exactly 
that, and `isNeededForSchema` is `false` in both the 3.5 and 4.x shims. So this 
helper splits any DV'd file bigger than `maxSplitBytes`, which by default is 
every file over 128 MB. The results are right. A single 177 KB DV'd file with 
small row groups, read with `spark.sql.files.maxPartitionBytes=4096`, came out 
as 44 native partitions and matched Spark. That works because DataFusion's 
`prune_by_range` skips row groups whose first page is outside the split, while 
the access plan stays in file coordinates.
   
   Could we add a test along those lines? It's the default shape for large 
files, and #5655 plans to change exactly this path. Could you also update 
#5655, which says splitting is disabled? Today every split fetches and decodes 
the whole DV, reads the footer, builds the whole-file plan and reserves memory 
for the whole file.



##########
contrib/delta-spark/src/test/scala/org/apache/comet/contrib/delta/CometDeltaNativeScanSuite.scala:
##########
@@ -0,0 +1,3584 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.comet.contrib.delta
+
+import java.io.File
+
+import scala.collection.mutable
+import scala.collection.mutable.ListBuffer
+import scala.concurrent.duration.DurationInt
+
+import org.apache.spark.scheduler.{SparkListener, SparkListenerTaskEnd}
+import org.apache.spark.sql.{DataFrame, Row}
+import org.apache.spark.sql.catalyst.expressions.{AttributeReference, 
DynamicPruningExpression, NamedExpression, StructsToJson}
+import org.apache.spark.sql.comet.CometDeltaNativeScanExec
+import org.apache.spark.sql.execution.{FileSourceScanExec, QueryExecution, 
ScalarSubquery, SparkPlan, SubqueryExec}
+import org.apache.spark.sql.execution.datasources.v2.V2TableWriteExec
+import org.apache.spark.sql.functions.{col, lit, to_json}
+import org.apache.spark.sql.internal.SQLConf
+import org.apache.spark.sql.types.{ByteType, LongType, StringType, 
StructField, StructType}
+import org.apache.spark.sql.util.QueryExecutionListener
+
+import org.apache.comet.CometConf
+import org.apache.comet.CometSparkSessionExtensions.isSpark40Plus
+import org.apache.comet.ExtendedExplainInfo
+import org.apache.comet.serde.OperatorOuterClass
+import org.apache.comet.serde.operator.CometNativeScan
+
+/**
+ * Differential suite: append-only Delta tables read through the native Delta 
scan must produce
+ * results identical to Spark's Delta reader, engage the native operator, and 
prune at row-group
+ * and page level.
+ */
+class CometDeltaNativeScanSuite extends CometDeltaTestBase {
+
+  test("plain delta table reads natively with identical results") {
+    withTempPath { dir =>
+      val path = dir.getAbsolutePath
+      spark
+        .range(0, 1000)
+        .selectExpr("id", "id * 2 as v", "cast(id as string) as s")
+        .write
+        .format("delta")
+        .save(path)
+
+      val df = spark.read.format("delta").load(path).filter(col("id") > 500)
+      checkDeltaNativeScanAnswer(df)
+    }
+  }
+
+  test("projection and filter on delta table") {
+    withTempPath { dir =>
+      val path = dir.getAbsolutePath
+      spark
+        .range(0, 1000)
+        .selectExpr("id", "id % 10 as bucket", "cast(id as double) as d")
+        .write
+        .format("delta")
+        .save(path)
+
+      val df = spark.read
+        .format("delta")
+        .load(path)
+        .select("bucket", "d")
+        .filter(col("d") < 100.0)
+      checkDeltaNativeScanAnswer(df)
+    }
+  }
+
+  test("partitioned delta table with partition filter") {
+    withTempPath { dir =>
+      val path = dir.getAbsolutePath
+      spark
+        .range(0, 1000)
+        .selectExpr("id", "id % 7 as p")
+        .write
+        .format("delta")
+        .partitionBy("p")
+        .save(path)
+
+      val df = spark.read.format("delta").load(path).filter(col("p") === 3)
+      checkDeltaNativeScanAnswer(df)
+      assert(df.count() > 0)
+    }
+  }
+
+  test("multi-file delta table after several appends") {
+    withTempPath { dir =>
+      val path = dir.getAbsolutePath
+      for (i <- 0 until 4) {
+        spark
+          .range(i * 100, (i + 1) * 100)
+          .selectExpr("id", "id * 3 as v")
+          .write
+          .format("delta")
+          .mode("append")
+          .save(path)
+      }
+      val df = spark.read.format("delta").load(path)
+      checkDeltaNativeScanAnswer(df)
+      assert(df.count() == 400)
+    }
+  }
+
+  test("time travel VERSION AS OF reads natively") {
+    withTempPath { dir =>
+      val path = dir.getAbsolutePath
+      spark.range(0, 100).write.format("delta").save(path)
+      spark.range(100, 200).write.format("delta").mode("append").save(path)
+
+      val v0 = spark.read.format("delta").option("versionAsOf", 0).load(path)
+      checkDeltaNativeScanAnswer(v0)
+      assert(v0.count() == 100)
+    }
+  }
+
+  test("selective predicate prunes row groups and pages") {
+    withTempPath { dir =>
+      val path = dir.getAbsolutePath
+      // Small row groups + page-level stats: sorted data so min/max stats are 
tight. The Delta
+      // writer ignores parquet.* DataFrameWriter options, so set them on the 
Hadoop conf.
+      val hadoopConf = spark.sparkContext.hadoopConfiguration
+      val oldBlockSize = hadoopConf.get("parquet.block.size")
+      val oldPageSize = hadoopConf.get("parquet.page.size")
+      hadoopConf.setInt("parquet.block.size", 256 * 1024)
+      hadoopConf.setInt("parquet.page.size", 16 * 1024)
+      try {
+        spark
+          .range(0, 500000)
+          .selectExpr("id", "id * 2 as v")
+          .sort("id")
+          .coalesce(1)
+          .write
+          .format("delta")
+          .save(path)
+      } finally {
+        if (oldBlockSize == null) hadoopConf.unset("parquet.block.size")
+        else hadoopConf.set("parquet.block.size", oldBlockSize)
+        if (oldPageSize == null) hadoopConf.unset("parquet.page.size")
+        else hadoopConf.set("parquet.page.size", oldPageSize)
+      }
+
+      def query = spark.read
+        .format("delta")
+        .load(path)
+        .filter(col("id") >= 100 && col("id") < 200)
+      checkDeltaNativeScanAnswer(query)
+
+      // checkSparkAnswer re-plans the query, so read metrics from a DataFrame 
we execute
+      // ourselves (collect() runs THIS Dataset's queryExecution; count() 
would plan a new one):
+      // its executed plan holds the metric objects native execution updated.
+      val df = query
+      assert(df.collect().length == 100)
+      val scans = deltaNativeScans(df)
+      assert(scans.size == 1)
+      val metrics = scans.head.metrics
+      val rowGroupsPruned = 
metrics.get("row_groups_pruned_statistics").map(_.value).getOrElse(0L)
+      val pagesPruned = 
metrics.get("page_index_rows_pruned").map(_.value).getOrElse(0L)
+      assert(
+        rowGroupsPruned > 0,
+        s"expected row-group pruning; metrics: ${metrics.map { case (k, v) => 
s"$k=${v.value}" }}")
+      assert(
+        pagesPruned > 0,
+        s"expected page-index pruning; metrics: ${metrics.map { case (k, v) =>
+            s"$k=${v.value}"
+          }}")
+    }
+  }
+
+  test("scalar subquery data filter is pushed down and prunes row groups and 
pages") {
+    withTempPath { dir =>
+      val path = s"${dir.getAbsolutePath}/data"
+      val thresholds = s"${dir.getAbsolutePath}/thresholds"
+      // Same layout as the selective-predicate test: small row groups + tight 
page stats.
+      val hadoopConf = spark.sparkContext.hadoopConfiguration
+      val oldBlockSize = hadoopConf.get("parquet.block.size")
+      val oldPageSize = hadoopConf.get("parquet.page.size")
+      hadoopConf.setInt("parquet.block.size", 256 * 1024)
+      hadoopConf.setInt("parquet.page.size", 16 * 1024)
+      try {
+        spark
+          .range(0, 500000)
+          .selectExpr("id", "id * 2 as v")
+          .sort("id")
+          .coalesce(1)
+          .write
+          .format("delta")
+          .save(path)
+      } finally {
+        if (oldBlockSize == null) hadoopConf.unset("parquet.block.size")
+        else hadoopConf.set("parquet.block.size", oldBlockSize)
+        if (oldPageSize == null) hadoopConf.unset("parquet.page.size")
+        else hadoopConf.set("parquet.page.size", oldPageSize)
+      }
+      spark
+        .sql("SELECT CAST(100 AS BIGINT) AS lo, CAST(200 AS BIGINT) AS hi")
+        .write
+        .format("delta")
+        .save(thresholds)
+
+      // Scalar subqueries are PlanExpressions: unresolved at planning, so the 
bounds can
+      // only reach the native reader via the execution-time 
resolve-and-append path.
+      def query = spark.sql(
+        s"SELECT * FROM delta.`$path` WHERE id >= (SELECT lo FROM 
delta.`$thresholds`) " +
+          s"AND id < (SELECT hi FROM delta.`$thresholds`)")
+      checkDeltaNativeScanAnswer(query)
+
+      val df = query
+      assert(df.collect().length == 100)
+      // The thresholds table inside the subquery is also claimed natively; 
pick the
+      // main data-table scan by its output.
+      assertSubqueryFilterPushed(df, dataColumn = "v")
+      val scans = deltaNativeScans(df).filter(_.output.exists(_.name == "v"))
+      assert(scans.size == 1)
+      val metrics = scans.head.metrics
+      val rowGroupsPruned = 
metrics.get("row_groups_pruned_statistics").map(_.value).getOrElse(0L)
+      val pagesPruned = 
metrics.get("page_index_rows_pruned").map(_.value).getOrElse(0L)
+      assert(
+        rowGroupsPruned > 0,
+        s"expected row-group pruning from the resolved subquery bounds; 
metrics: ${metrics.map {
+            case (k, v) => s"$k=${v.value}"
+          }}")
+      assert(
+        pagesPruned > 0,
+        s"expected page-index pruning from the resolved subquery bounds; 
metrics: ${metrics.map {
+            case (k, v) => s"$k=${v.value}"
+          }}")
+    }
+  }
+
+  test("deletion vectors: scalar subquery filter composes with DV 
application") {
+    withTempPath { dir =>
+      val path = s"${dir.getAbsolutePath}/data"
+      val thresholds = s"${dir.getAbsolutePath}/thresholds"
+      createDvTable(path, rows = 10000)
+      spark.sql(s"DELETE FROM delta.`$path` WHERE id % 2 = 0")
+      spark
+        .sql("SELECT CAST(5000 AS BIGINT) AS lo")
+        .write
+        .format("delta")
+        .save(thresholds)
+
+      def query =
+        spark.sql(s"SELECT * FROM delta.`$path` WHERE id >= (SELECT lo FROM 
delta.`$thresholds`)")
+      checkDeltaNativeScanAnswer(query)
+      // Deleted rows must stay deleted with the pushed bound applied in-scan.
+      val df = query
+      val rows = df.collect()
+      assert(rows.length == 2500)
+      assert(rows.forall(r => r.getLong(0) % 2 == 1 && r.getLong(0) >= 5000))
+      assertSubqueryFilterPushed(df, dataColumn = "v")
+    }
+  }
+
+  test("column mapping: scalar subquery filter on a renamed column") {
+    withTempPath { dir =>
+      val path = s"${dir.getAbsolutePath}/data"
+      val thresholds = s"${dir.getAbsolutePath}/thresholds"
+      spark.range(0, 1000).selectExpr("id", "id * 2 as 
v").write.format("delta").save(path)
+      enableColumnMapping(path)
+      spark.sql(s"ALTER TABLE delta.`$path` RENAME COLUMN v TO w")
+      spark
+        .sql("SELECT CAST(900 AS BIGINT) AS lo")
+        .write
+        .format("delta")
+        .save(thresholds)
+
+      // The pushed filter references the renamed column: it must bind against 
the
+      // physical read schema, not the logical name.
+      def query =
+        spark.sql(s"SELECT * FROM delta.`$path` WHERE w >= (SELECT lo FROM 
delta.`$thresholds`)")
+      checkDeltaNativeScanAnswer(query)
+      val df = query
+      assert(df.collect().length == 550)
+      assertSubqueryFilterPushed(df, dataColumn = "w")
+    }
+  }
+
+  /**
+   * Assert the resolved scalar-subquery bound was actually appended to the 
native scan's
+   * execution-time common data (answers alone cannot show this: Spark's 
covering FilterExec would
+   * mask a silently-skipped pushdown). `df` must already have been executed.
+   */
+  private def assertSubqueryFilterPushed(df: DataFrame, dataColumn: String): 
Unit = {
+    val scans = deltaNativeScans(df).collect {
+      case s: CometDeltaNativeScanExec if s.output.exists(_.name == 
dataColumn) => s
+    }
+    assert(scans.size == 1)
+    val scan = scans.head
+    val planTimeFilters =
+      
DeltaSparkScanEnvelope.unpack(scan.nativeOp).getCommon.getDataFiltersCount
+    val executedFilters = OperatorOuterClass.DeltaSparkScan
+      .parseFrom(scan.commonData)
+      .getCommon
+      .getDataFiltersCount
+    assert(
+      executedFilters > planTimeFilters,
+      "expected resolved subquery filters appended at execution: " +
+        s"plan-time=$planTimeFilters executed=$executedFilters " +
+        s"dataFilters=${scan.dataFilters.mkString("; ")}")
+  }
+
+  test("scalar subquery filter is NOT pushed below a limit") {
+    withTempPath { dir =>
+      val path = dir.getAbsolutePath
+      spark.range(0, 3).selectExpr("id").write.format("delta").save(path)
+      
spark.read.format("delta").load(path).createOrReplaceTempView("t_limit_pushdown")
+
+      val df = spark.sql(
+        "SELECT id FROM (SELECT id FROM t_limit_pushdown ORDER BY id LIMIT 1) 
q " +
+          "WHERE id > (SELECT max(id) FROM range(1))")
+      checkSparkAnswer(df)
+      assert(df.collect().isEmpty)
+      assertNoSubqueryFilterPushed(df)
+    }
+  }
+
+  test("scalar subquery filter is NOT pushed across a nondeterministic 
projection") {
+    withSQLConf(CometConf.COMET_PARQUET_ROW_FILTER_PUSHDOWN_ENABLED.key -> 
"true") {
+      withTempPath { dir =>
+        val path = dir.getAbsolutePath
+        spark.range(0, 5).coalesce(1).write.format("delta").save(path)
+        
spark.read.format("delta").load(path).createOrReplaceTempView("t_monotonic_id")
+
+        // A deterministic conjunct does not commute with a nondeterministic 
projection: the
+        // subquery bound must not be pushed into the scan below `seq`, or the 
surviving rows'
+        // monotonically_increasing_id() values change and the answer is wrong.
+        val df = spark.sql(
+          "SELECT id FROM (SELECT id, monotonically_increasing_id() AS seq " +
+            "FROM t_monotonic_id) q WHERE id > (SELECT max(id) FROM range(1)) 
AND seq = 1")
+        checkSparkAnswer(df)
+        assert(df.collect().toSeq == Seq(Row(1)))
+        assertNoSubqueryFilterPushed(df)
+      }
+    }
+  }
+
+  /**
+   * Assert no scalar-subquery filter was harvested and pushed into the native 
scan's
+   * execution-time common data: the scan must sit below a non-commuting 
operator (e.g. LIMIT /
+   * TopN), so the covering FilterExec's predicate must stay above it rather 
than move into the
+   * scan. Also confirms the query still engaged the native Delta scan, i.e. 
this exercises the
+   * commutativity guard rather than a plan that fell back to Spark entirely. 
`df` must already
+   * have been executed.
+   */
+  private def assertNoSubqueryFilterPushed(df: DataFrame): Unit = {
+    val scans = deltaNativeScans(df).collect { case s: 
CometDeltaNativeScanExec => s }
+    assert(scans.size == 1, s"expected exactly one native Delta scan; found 
${scans.size}")
+    val scan = scans.head
+    val planTimeFilters =
+      
DeltaSparkScanEnvelope.unpack(scan.nativeOp).getCommon.getDataFiltersCount
+    val executedFilters = OperatorOuterClass.DeltaSparkScan
+      .parseFrom(scan.commonData)
+      .getCommon
+      .getDataFiltersCount
+    assert(
+      executedFilters == planTimeFilters,
+      "expected no subquery filter pushed across the non-commuting operator 
between the " +
+        s"covering filter and the scan: plan-time=$planTimeFilters 
executed=$executedFilters " +
+        s"dataFilters=${scan.dataFilters.mkString("; ")}")
+  }
+
+  test("scalar subquery filter rejected by serde still marks the scan as 
filtered") {
+    withTempPath { dir =>
+      val path = s"${dir.getAbsolutePath}/data"
+      val bounds = s"${dir.getAbsolutePath}/bounds"
+      spark.range(0, 100).selectExpr("id", "id * 2 as 
v").write.format("delta").save(path)
+      spark.sql("SELECT CAST(42 AS BIGINT) AS 
lo").write.format("delta").save(bounds)
+
+      // With EqualNullSafe disabled the resolved bound cannot serialize, yet 
the scan must still
+      // carry has_data_filters so native treats it as a filtered read, 
exactly like core does.
+      withSQLConf("spark.comet.expression.EqualNullSafe.enabled" -> "false") {
+        def query =
+          spark.sql(
+            s"SELECT * FROM delta.`$path` WHERE id <=> (SELECT max(lo) FROM 
delta.`$bounds`)")
+        checkDeltaNativeScanAnswer(query)
+        val df = query
+        assert(df.collect().toSeq == Seq(Row(42L, 84L)))
+        assertUnserializedSubqueryFilterMarksScanFiltered(df, dataColumn = "v")
+      }
+    }
+  }
+
+  test("unserializable scalar subquery filter keeps the safe TIMESTAMP_MILLIS 
conversion") {
+    // Same fixture as core's "filtered TIMESTAMP_MILLIS scans do not convert 
values Spark can
+    // skip": a raw file whose only overflowing millisecond value Spark prunes 
from the footer
+    // statistics once the resolved bound is pushed, so native must not 
convert it either.
+    withTempPath { dir =>
+      val path = s"${dir.getAbsolutePath}/data"
+      val bounds = s"${dir.getAbsolutePath}/bounds"
+      writeRawParquetFile(
+        path,
+        """message root {
+          |  optional int32 id;
+          |  optional int64 ts(TIMESTAMP_MILLIS);
+          |}""".stripMargin) { factory =>
+        (1 to 16).map(id => factory.newGroup().append("id", id).append("ts", 
1717243200000L)) :+
+          factory.newGroup().append("id", 17).append("ts", 9223372036854776L)
+      }
+      spark.sql(s"CONVERT TO DELTA parquet.`$path` NO STATISTICS")
+      spark.sql("SELECT timestamp_seconds(0) AS 
bound").write.format("delta").save(bounds)
+
+      withSQLConf(
+        "spark.comet.expression.EqualNullSafe.enabled" -> "false",
+        "spark.sql.parquet.datetimeRebaseModeInRead" -> "CORRECTED",
+        "spark.sql.parquet.int96RebaseModeInRead" -> "CORRECTED") {
+        def query = spark.sql(
+          s"SELECT id, ts FROM delta.`$path` " +
+            s"WHERE ts <=> (SELECT max(bound) FROM delta.`$bounds`)")
+        // Spark 3.x never pushes subquery filters into its parquet reader and 
converts the
+        // overflowing value itself, so the answer comparison is meaningful on 
Spark 4.0+ only.
+        if (isSpark40Plus) {
+          checkDeltaNativeScanAnswer(query)
+        }
+        val df = query
+        assert(df.collect().isEmpty)
+        assert(
+          deltaNativeScans(df).nonEmpty,
+          s"expected a native Delta scan:\n${df.queryExecution}")
+        assertUnserializedSubqueryFilterMarksScanFiltered(df, dataColumn = 
"ts")
+      }
+    }
+  }
+
+  /**
+   * Assert the execution-time common data of the scan producing `dataColumn` 
reports
+   * `has_data_filters` with no serialized data filter: the plan-time proto 
carries neither, and
+   * the resolved subquery filter is the only data filter, so only the 
execution-time path can set
+   * the bit. `df` must already have been executed.
+   */
+  private def assertUnserializedSubqueryFilterMarksScanFiltered(
+      df: DataFrame,
+      dataColumn: String): Unit = {
+    val scans = deltaNativeScans(df).collect {
+      case s: CometDeltaNativeScanExec if s.output.exists(_.name == 
dataColumn) => s
+    }
+    assert(scans.size == 1, s"expected exactly one native Delta scan; found 
${scans.size}")
+    val scan = scans.head
+    assert(
+      scan.dataFilters.exists(_.exists(_.isInstanceOf[ScalarSubquery])),
+      s"expected a scalar subquery data filter: ${scan.dataFilters.mkString("; 
")}")
+    val planTime = DeltaSparkScanEnvelope.unpack(scan.nativeOp).getCommon
+    assert(!planTime.getHasDataFilters && planTime.getDataFiltersCount == 0)
+    val executed = 
OperatorOuterClass.DeltaSparkScan.parseFrom(scan.commonData).getCommon
+    assert(
+      executed.getHasDataFilters,
+      "expected has_data_filters at execution even though the resolved 
subquery filter did " +
+        s"not serialize: dataFilters=${scan.dataFilters.mkString("; ")}")
+    assert(executed.getDataFiltersCount == 0)
+  }
+
+  test("aggregation over delta table") {
+    withTempPath { dir =>
+      val path = dir.getAbsolutePath
+      spark
+        .range(0, 10000)
+        .selectExpr("id", "id % 13 as g", "id * 2 as v")
+        .write
+        .format("delta")
+        .save(path)
+
+      val df = spark.read
+        .format("delta")
+        .load(path)
+        .groupBy("g")
+        .sum("v")
+      checkDeltaNativeScanAnswer(df)
+    }
+  }
+
+  test("conf disables the native delta scan") {
+    withTempPath { dir =>
+      val path = dir.getAbsolutePath
+      spark.range(0, 100).write.format("delta").save(path)
+
+      withSQLConf(DeltaScanConf.COMET_DELTA_NATIVE_ENABLED.key -> "false") {
+        val df = spark.read.format("delta").load(path)
+        checkSparkAnswer(df)
+        assert(deltaNativeScans(df).isEmpty)
+      }
+    }
+  }
+
+  test("native delta scan is opt-in: disabled when the conf is not set") {
+    withTempPath { dir =>
+      val path = dir.getAbsolutePath
+      spark.range(0, 100).write.format("delta").save(path)
+
+      // The suite base enables the scan globally; drop the key entirely to
+      // observe the out-of-the-box default.
+      spark.conf.unset(DeltaScanConf.COMET_DELTA_NATIVE_ENABLED.key)
+      try {
+        assert(!DeltaScanConf.scanEnabled)
+        val df = spark.read.format("delta").load(path)
+        checkSparkAnswer(df)
+        assert(deltaNativeScans(df).isEmpty)
+      } finally {
+        spark.conf.set(DeltaScanConf.COMET_DELTA_NATIVE_ENABLED.key, "true")
+      }
+    }
+  }
+
+  test("table root under a directory whose name contains a newline falls back 
to Spark") {
+    // object_store recognizes the `file` scheme but rejects the control 
character in the
+    // directory name (`%0A` in the URI), so native execution could not open 
the table where
+    // Spark's Hadoop-backed reader can. The claim gate must decline before 
native planning.
+    withTempPath { dir =>
+      val path = new File(new File(dir, "dir\n"), "data").getAbsolutePath
+      spark.range(0, 100).write.format("delta").save(path)
+
+      val df = spark.read.format("delta").load(path)
+      assert(
+        deltaNativeScans(df).isEmpty,
+        "Expected no native Delta scan under a newline directory:\n" +
+          s"${df.queryExecution.executedPlan}")
+      checkSparkAnswerAndFallbackReason(
+        df,
+        "Native Delta scan cannot open path 'file:" + dir.getAbsolutePath +
+          "/dir%0A/data': object_store rejects it")
+    }
+  }
+
+  test("shallow clone whose source data files sit under a newline directory 
falls back") {
+    // The clone's own root is an ordinary path, so only the selected data 
files (resolved to
+    // the source table's directory) carry the rejected segment: this 
exercises the
+    // selected-paths probe, not the root gate. The reason names the first 
such complete path,
+    // a data file under the source directory.
+    withTempPath { dir =>
+      val sourcePath = new File(new File(dir, "dir\n"), 
"source").getAbsolutePath
+      val clonePath = new File(dir, "clone").getAbsolutePath
+      spark.range(0, 100).write.format("delta").save(sourcePath)
+      spark.sql(s"CREATE TABLE delta.`$clonePath` SHALLOW CLONE 
delta.`$sourcePath`")
+
+      val df = spark.read.format("delta").load(clonePath)
+      assert(
+        deltaNativeScans(df).isEmpty,
+        "Expected no native Delta scan for a clone of a newline-directory 
source:\n" +
+          s"${df.queryExecution.executedPlan}")
+      checkSparkAnswerAndFallbackReason(
+        df,
+        "Native Delta scan cannot open path 'file:" + dir.getAbsolutePath +
+          "/dir%0A/source/")
+    }
+  }
+
+  test("converted Parquet table with a newline in a data file basename falls 
back to Spark") {
+    // CONVERT TO DELTA keeps the existing Parquet file names, so the rejected 
character sits in
+    // the basename rather than a directory segment: the table root and every 
parent directory
+    // pass the path probe, and only a check of the complete selected path can 
decline.
+    withTempPath { dir =>
+      val path = new File(dir, "data").getAbsolutePath
+      spark.range(0, 100).repartition(2).write.parquet(path)
+      val original = new 
File(path).listFiles().filter(_.getName.endsWith(".parquet")).head
+      val renamed = new File(path, "part-00000\n.snappy.parquet")
+      java.nio.file.Files.move(original.toPath, renamed.toPath)
+      spark.sql(s"CONVERT TO DELTA parquet.`$path`")
+
+      val df = spark.read.format("delta").load(path)
+      assert(
+        deltaNativeScans(df).isEmpty,
+        "Expected no native Delta scan for a converted table with a newline 
basename:\n" +
+          s"${df.queryExecution.executedPlan}")
+      checkSparkAnswerAndFallbackReason(
+        df,
+        s"Native Delta scan cannot open path 
'file:$path/part-00000%0A.snappy.parquet': " +
+          "object_store rejects it")
+    }
+  }
+
+  private def createDvTable(path: String, rows: Long = 1000): Unit = {
+    spark.range(0, rows).selectExpr("id", "id * 2 as 
v").write.format("delta").save(path)
+    spark.sql(
+      s"ALTER TABLE delta.`$path` SET TBLPROPERTIES 
('delta.enableDeletionVectors' = 'true')")
+  }
+
+  /**
+   * Same shape as `createDvTable`, plus one extra TINYINT column (value 7) 
under `columnName`.
+   */
+  private def createDvTableWithExtraColumn(
+      path: String,
+      columnName: String,
+      rows: Long = 1000): Unit = {
+    spark
+      .range(0, rows)
+      .selectExpr("id", s"cast(7 as tinyint) as `$columnName`")
+      .write
+      .format("delta")
+      .save(path)
+    spark.sql(
+      s"ALTER TABLE delta.`$path` SET TBLPROPERTIES 
('delta.enableDeletionVectors' = 'true')")
+  }
+
+  test("deletion vectors: DELETE-produced DVs read natively with correct 
results") {
+    withTempPath { dir =>
+      val path = dir.getAbsolutePath
+      createDvTable(path)
+      spark.sql(s"DELETE FROM delta.`$path` WHERE id % 2 = 0")
+
+      val df = spark.read.format("delta").load(path)
+      checkDeltaNativeScanAnswer(df)
+      assert(df.count() == 500)
+    }
+  }
+
+  test(
+    "deletion vectors: user column named like the synthetic internal-column 
slot keeps its " +
+      "own values") {
+    withTempPath { dir =>
+      val path = dir.getAbsolutePath
+      val collidingName = "_comet_delta___delta_internal_is_row_deleted"
+      createDvTableWithExtraColumn(path, collidingName)
+      spark.sql(s"DELETE FROM delta.`$path` WHERE id = 0")
+
+      val df = spark.read.format("delta").load(path).select("id", 
collidingName)
+      checkDeltaNativeScanAnswer(df)
+      val survivingValues = 
df.collect().map(_.getAs[Byte](collidingName)).distinct
+      assert(
+        survivingValues.sameElements(Array(7.toByte)),
+        "expected the user column's own value (7) to survive DV filtering, " +
+          s"got ${survivingValues.toSeq}")
+    }
+  }
+
+  test("deletion vectors: normally named extra column alongside DVs reads 
natively") {
+    withTempPath { dir =>
+      val path = dir.getAbsolutePath
+      createDvTableWithExtraColumn(path, "tag")
+      spark.sql(s"DELETE FROM delta.`$path` WHERE id = 0")
+
+      val df = spark.read.format("delta").load(path).select("id", "tag")
+      checkDeltaNativeScanAnswer(df)
+      val survivingValues = df.collect().map(_.getAs[Byte]("tag")).distinct
+      assert(
+        survivingValues.sameElements(Array(7.toByte)),
+        "expected the extra column's value (7) to survive DV filtering, got " +
+          survivingValues.toSeq)
+    }
+  }
+
+  test("deletion vectors: UPDATE-produced DVs read natively with correct 
results") {
+    withTempPath { dir =>
+      val path = dir.getAbsolutePath
+      createDvTable(path)
+      spark.sql(s"UPDATE delta.`$path` SET v = -1 WHERE id < 100")
+
+      val df = spark.read.format("delta").load(path)
+      checkDeltaNativeScanAnswer(df)
+      assert(df.filter(col("v") === -1).count() == 100)
+      assert(df.count() == 1000)
+    }
+  }
+
+  test("deletion vectors: multiple DELETEs accumulate correctly") {
+    withTempPath { dir =>
+      val path = dir.getAbsolutePath
+      createDvTable(path)
+      spark.sql(s"DELETE FROM delta.`$path` WHERE id % 2 = 0")
+      spark.sql(s"DELETE FROM delta.`$path` WHERE id % 3 = 0")
+
+      val df = spark.read.format("delta").load(path)
+      checkDeltaNativeScanAnswer(df)
+      // odd ids not divisible by 3
+      assert(df.count() == (0L until 1000L).count(i => i % 2 != 0 && i % 3 != 
0))
+    }
+  }
+
+  test(
+    "deletion vectors: maxDeletedRowsPerFile budget declines an oversized DV 
and " +
+      "claims once raised") {
+    withTempPath { dir =>
+      val path = dir.getAbsolutePath
+      // repartition(4) guarantees >= 2 physical files so the per-file 
cardinality gate has
+      // more than one file to inspect, mirroring design F3's multi-file test 
shape.
+      spark
+        .range(0, 1000)
+        .selectExpr("id", "id * 2 as v")
+        .repartition(4)
+        .write
+        .format("delta")
+        .save(path)
+      spark.sql(
+        s"ALTER TABLE delta.`$path` SET TBLPROPERTIES 
('delta.enableDeletionVectors' = 'true')")
+      spark.sql(s"DELETE FROM delta.`$path` WHERE id % 2 = 0")
+
+      withSQLConf(DeltaScanConf.COMET_DELTA_MAX_DELETED_ROWS_PER_FILE.key -> 
"1") {
+        val df = spark.read.format("delta").load(path)
+        checkSparkAnswer(df)
+        assert(
+          deltaNativeScans(df).isEmpty,
+          "a budget of 1 deleted row per file must decline every DV-bearing 
file")
+      }
+
+      withSQLConf(DeltaScanConf.COMET_DELTA_MAX_DELETED_ROWS_PER_FILE.key -> 
"1000000") {
+        val df = spark.read.format("delta").load(path)
+        checkDeltaNativeScanAnswer(df)
+      }
+    }
+  }
+
+  test("deletion vectors: maxDeletedRowsPerFile decline reason names the conf 
key") {
+    withTempPath { dir =>
+      val path = dir.getAbsolutePath
+      createDvTable(path)
+      spark.sql(s"DELETE FROM delta.`$path` WHERE id % 2 = 0")
+
+      withSQLConf(DeltaScanConf.COMET_DELTA_MAX_DELETED_ROWS_PER_FILE.key -> 
"1") {
+        checkSparkAnswerAndFallbackReason(
+          spark.read.format("delta").load(path),
+          DeltaScanConf.COMET_DELTA_MAX_DELETED_ROWS_PER_FILE.key)
+      }
+    }
+  }
+
+  test("deletion vectors: fully-deleted region and selective predicate still 
prune pages") {
+    withTempPath { dir =>
+      val path = dir.getAbsolutePath
+      val hadoopConf = spark.sparkContext.hadoopConfiguration
+      val oldBlockSize = hadoopConf.get("parquet.block.size")
+      val oldPageSize = hadoopConf.get("parquet.page.size")
+      hadoopConf.setInt("parquet.block.size", 256 * 1024)
+      hadoopConf.setInt("parquet.page.size", 16 * 1024)
+      try {
+        spark
+          .range(0, 500000)
+          .selectExpr("id", "id * 2 as v")
+          .sort("id")
+          .coalesce(1)
+          .write
+          .format("delta")
+          .save(path)
+      } finally {
+        if (oldBlockSize == null) hadoopConf.unset("parquet.block.size")
+        else hadoopConf.set("parquet.block.size", oldBlockSize)
+        if (oldPageSize == null) hadoopConf.unset("parquet.page.size")
+        else hadoopConf.set("parquet.page.size", oldPageSize)
+      }
+      spark.sql(
+        s"ALTER TABLE delta.`$path` SET TBLPROPERTIES 
('delta.enableDeletionVectors' = 'true')")
+      // Delete a slice inside the predicate range and a large slice outside 
it.
+      spark.sql(s"DELETE FROM delta.`$path` WHERE id >= 150 AND id < 160")
+      spark.sql(s"DELETE FROM delta.`$path` WHERE id >= 300000")
+
+      def query = spark.read
+        .format("delta")
+        .load(path)
+        .filter(col("id") >= 100 && col("id") < 200)
+      checkDeltaNativeScanAnswer(query)
+
+      val df = query
+      assert(df.collect().length == 90)
+      val scans = deltaNativeScans(df)
+      assert(scans.size == 1)
+      val metrics = scans.head.metrics
+      val pagesPruned = 
metrics.get("page_index_rows_pruned").map(_.value).getOrElse(0L)
+      assert(
+        pagesPruned > 0,
+        s"expected page-index pruning to compose with DVs; metrics: 
${metrics.map { case (k, v) =>
+            s"$k=${v.value}"
+          }}")
+    }
+  }
+
+  test("deletion vectors: aggregation over DV table") {
+    withTempPath { dir =>
+      val path = dir.getAbsolutePath
+      createDvTable(path, rows = 10000)
+      spark.sql(s"DELETE FROM delta.`$path` WHERE id % 7 = 0")
+
+      val df = spark.read.format("delta").load(path).groupBy(col("id") % 
13).count()
+      checkDeltaNativeScanAnswer(df)
+    }
+  }
+
+  test("deletion vectors: partitioned table reads natively with correct 
results") {
+    withTempPath { dir =>
+      val path = dir.getAbsolutePath
+      spark
+        .range(0, 1000)
+        .selectExpr("id", "id % 5 as p", "id * 2 as v")
+        .write
+        .format("delta")
+        .partitionBy("p")
+        .save(path)
+      spark.sql(
+        s"ALTER TABLE delta.`$path` SET TBLPROPERTIES 
('delta.enableDeletionVectors' = 'true')")
+      spark.sql(s"DELETE FROM delta.`$path` WHERE id % 3 = 0")
+
+      val df = spark.read.format("delta").load(path).filter(col("p") === 2)
+      checkDeltaNativeScanAnswer(df)
+      assert(df.count() == (0L until 1000L).count(i => i % 5 == 2 && i % 3 != 
0))
+
+      val all = spark.read.format("delta").load(path)
+      checkDeltaNativeScanAnswer(all)
+      assert(all.count() == (0L until 1000L).count(_ % 3 != 0))
+    }
+  }
+
+  test("deletion vectors: combined with constant metadata columns") {
+    withTempPath { dir =>
+      val path = dir.getAbsolutePath
+      createDvTable(path)
+      spark.sql(s"DELETE FROM delta.`$path` WHERE id < 250")
+
+      val df = spark.read
+        .format("delta")
+        .load(path)
+        .selectExpr("id", "v", "_metadata.file_name as fn")
+      checkSparkAnswer(df.selectExpr("id", "v", "length(fn) > 0"))
+      // Whether this claims or declines, results must match; if it claimed, 
verify the
+      // native node is present so the combination is actually exercised when 
supported.
+      val rows = df.collect()
+      assert(rows.length == 750)
+      assert(rows.forall(_.getString(2).nonEmpty))
+    }
+  }
+
+  test(
+    "deletion vectors: constant-metadata field names are deduplicated against 
the physical " +
+      "data and partition schemas") {
+    // End-to-end coverage is not possible here: selecting any `_metadata.*` 
field in the DV
+    // shape always declines today for an unrelated, pre-existing reason -- 
Spark reuses the
+    // scan's own row-index bookkeeping attribute as `_metadata.row_index`'s 
source, and
+    // `DeltaScanSupport.rowIndexUnusedAbove` conservatively treats extracting 
ANY `_metadata`
+    // field as making that attribute live (see "combined with constant 
metadata columns"
+    // above, which hedges its assertions for the same reason). That decline 
fires before
+    // `buildDvScanCommon` ever runs, regardless of collision, so it cannot 
exercise the fix.
+    // Test the builder's dedup logic directly instead, the same way 
`storeUris` and
+    // `mergedObjectStoreOptions` are unit-tested without a live scan.
+    val physicalDataSchema =
+      StructType(Seq(StructField("_comet_metadata_file_path", ByteType)))
+    val physicalPartitionSchema =
+      StructType(Seq(StructField("_comet_metadata_file_size", LongType)))
+    val fileConstantMetadataColumns = Seq(
+      AttributeReference("file_path", StringType, nullable = false)(),
+      AttributeReference("file_size", LongType, nullable = false)())
+
+    val constantMetadataFields = CometNativeScan.uniqueConstantMetadataFields(
+      fileConstantMetadataColumns,
+      physicalDataSchema.fields.map(_.name).toSet ++ 
physicalPartitionSchema.fields
+        .map(_.name)
+        .toSet)
+    assert(
+      constantMetadataFields.map(_.name) == Seq(
+        "_comet_metadata_file_path_",
+        "_comet_metadata_file_size_"),
+      "expected both constant-metadata names to be uniquified on collision, 
got " +
+        s"${constantMetadataFields.map(_.name)}")
+
+    // The DV builder must feed these already-unique names into 
allocateUniqueInternalFields's
+    // reserved set so the internal-column suffix chain stays consistent with 
them.
+    val requiredSchema = StructType(
+      Seq(
+        StructField("id", LongType),
+        StructField(CometDeltaNativeScan.IsRowDeletedColumn, ByteType),
+        StructField(CometDeltaNativeScan.RowIndexColumn, LongType)))
+    val internalFields = CometDeltaNativeScan.allocateUniqueInternalFields(
+      requiredSchema,
+      physicalDataSchema,
+      physicalPartitionSchema,
+      constantMetadataFields)
+
+    val allNames = physicalDataSchema.fields.map(_.name) ++
+      physicalPartitionSchema.fields.map(_.name) ++
+      constantMetadataFields.map(_.name) ++
+      internalFields.map(_.name)
+    assert(allNames.distinct.length == allNames.length, s"expected all names 
distinct: $allNames")
+  }
+
+  test(
+    "non-DV shape: user column named like the synthetic constant-metadata slot 
keeps its " +
+      "own values") {
+    withTempPath { dir =>
+      val path = dir.getAbsolutePath
+      val collidingName = "_comet_metadata_file_path"
+      spark
+        .range(0, 100)
+        .selectExpr("id", s"cast(7 as tinyint) as `$collidingName`")
+        .write
+        .format("delta")
+        .save(path)
+
+      val df = spark.read
+        .format("delta")
+        .load(path)
+        .selectExpr("id", s"`$collidingName`", "_metadata.file_path as fp")
+      checkDeltaNativeScanAnswer(df)
+      val rows = df.collect()
+      val survivingValues = rows.map(_.getAs[Byte](collidingName)).distinct
+      assert(
+        survivingValues.sameElements(Array(7.toByte)),
+        "expected the user column's own value (7) to survive the 
constant-metadata " +
+          s"collision, got ${survivingValues.toSeq}")
+      assert(
+        rows.forall(_.getString(2).nonEmpty),
+        "expected _metadata.file_path to still report a real path")
+    }
+  }
+
+  test("deletion vectors: special characters in table path") {
+    withTempDir { base =>
+      val dir = new java.io.File(base, "s p a r k %dv% test")
+      val path = dir.getAbsolutePath
+      createDvTable(path)
+      spark.sql(s"DELETE FROM delta.`$path` WHERE id % 2 = 0")
+
+      val df = spark.read.format("delta").load(path)
+      checkDeltaNativeScanAnswer(df)
+      assert(df.count() == 500)
+    }
+  }
+
+  test("deletion vectors: decline when row_index is consumed via multi-hop 
aliases") {
+    withTempPath { dir =>
+      val path = dir.getAbsolutePath
+      createDvTable(path)
+      spark.sql(s"DELETE FROM delta.`$path` WHERE id % 2 = 0")
+
+      val df = spark.read
+        .format("delta")
+        .load(path)
+        .selectExpr("id", "_metadata.row_index as ri")
+        .selectExpr("id", "ri + 1 as ri2")
+        .filter(col("ri2") > 10)
+      checkSparkAnswer(df)
+      assert(deltaNativeScans(df).isEmpty, "derived row_index consumption must 
decline")
+    }
+  }
+
+  test("deletion vectors: decline when row_index feeds a non-Project 
operator") {
+    withTempPath { dir =>
+      val path = dir.getAbsolutePath
+      createDvTable(path)
+      spark.sql(s"DELETE FROM delta.`$path` WHERE id % 2 = 0")
+
+      val df = spark.read
+        .format("delta")
+        .load(path)
+        .groupBy(col("_metadata.row_index") % 7)
+        .count()
+      checkSparkAnswer(df)
+      assert(deltaNativeScans(df).isEmpty, "aggregate over row_index must 
decline")
+    }
+  }
+
+  test("deletion vectors: decline when _metadata.row_index is referenced above 
the scan") {
+    withTempPath { dir =>
+      val path = dir.getAbsolutePath
+      createDvTable(path)
+      spark.sql(s"DELETE FROM delta.`$path` WHERE id % 2 = 0")
+
+      val df = spark.read
+        .format("delta")
+        .load(path)
+        .selectExpr("id", "_metadata.row_index as ri")
+      checkSparkAnswer(df)
+      assert(
+        deltaNativeScans(df).isEmpty,
+        "plans consuming a real row_index must fall back to Spark")
+    }
+  }
+
+  /**
+   * Every [[SparkPlan]] executed during `body`, captured via a 
[[QueryExecutionListener]] rather
+   * than a returned `DataFrame`'s own plan: a `DataFrameWriter` action such 
as `.write.parquet`
+   * has no result `Dataset` to call `.queryExecution` on, so the write's 
physical plan -- the one
+   * `DeltaScanSupport.declineReason` actually saw -- is only observable this 
way.
+   */
+  private def capturePlansDuring(body: => Unit): Seq[SparkPlan] = {

Review Comment:
   This test failed in one of my two full local runs on Spark 3.5, with an 
empty reason list. It passes on its own. `QueryExecutionListener` callbacks 
arrive on the listener bus asynchronously, and this helper unregisters as soon 
as `body` returns. So the write's plan can arrive after we've stopped 
listening. When that happens `nativeScans.isEmpty` passes vacuously and only 
the reason check notices. A 300 ms sleep in `onSuccess` makes it fail every 
time. Adding `CometListenerBusUtils.waitUntilEmpty(spark.sparkContext)` after 
`body` makes it pass again with the sleep still there, which is what 
`CometIcebergTestBase.capturePlans` does.
   
   Could both copies of this helper drain the bus, this one and the one in 
`CometDeltaDmlReproSuite`? The comment on `collectTaskInputMetrics` says the 
suite can't reach `waitUntilEmpty`, but 
`org.apache.spark.CometListenerBusUtils` is in the spark test-jar this module 
already depends on. `delta_3_5` runs in the merge queue for almost any change 
under `native/` or `spark/src/main`, so a flake here would evict other people's 
PRs.



##########
pom.xml:
##########
@@ -45,11 +45,10 @@ under the License.
     <project.reporting.outputEncoding>UTF-8</project.reporting.outputEncoding>
     <java.version>17</java.version>
     <!--
-      Default Delta dep version (only consumed when `-Pcontrib-delta` is also
-      active). Spark profiles override this: spark-3.5 -> 3.3.2, spark-4.1 ->
-      4.1.0. The top-level default lets Maven invocations that don't activate a
-      Spark profile (e.g. `mvn -Pcontrib-delta spotless:apply`) resolve the
-      property without an error.
+      Default Delta dep version, read by both Delta contribs (`-Pcontrib-delta`

Review Comment:
   Following up on my earlier `delta.version` comment. This comment says the 
property is read by both contribs and that every Spark profile overrides it. 
But `spark/pom.xml` overrides it again in its own Spark profiles for 
`-Pcontrib-delta`, so the two modules resolve different versions on the same 
profile. `help:evaluate` gives 4.0.0 for `spark` and 4.0.1 for 
`contrib/delta-spark` on `spark-4.0`, and 4.1.0 against 4.3.1 on `spark-4.1`. 
Could this module use its own property, something like `delta.spark.version`, 
so bumping one pairing can't quietly leave the other behind?



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to