use std::{error::Error as _, fmt};
use datafusion::arrow::datatypes::SchemaRef;
use delta_arrow_reader::{
DeltaReaderError, DeltaReaderPhase, DeltaSnapshotSelection, DeltaTableBuilder,
DeltaTableSnapshot, datafusion::ScanOptions,
};
use crate::{
DeltaFunnelError, DeltaProtocolReport, DeltaSourceConfig, PhaseTimingReport,
RegisteredDeltaSource, observability,
progress::{ProgressEvent, ProgressOperation, ProgressPhase, ProgressReporter},
query_engine::datafusion::{
register_delta_source_with_scan_options, reject_existing_delta_registration_name,
validate_delta_table_snapshot_protocol,
},
report::PhaseTimer,
support::{sanitize_text_for_display, sanitize_uri_for_display},
table_formats::{S3AuthModeHint, validate_table_source_names},
};
use super::super::{DeltaFunnelSession, LazyTable};
const SOURCE_LOADING_PHASE: &str = "source_loading";
const PROTOCOL_PREFLIGHT_PHASE: &str = "protocol_preflight";
const DATAFUSION_REGISTRATION_PHASE: &str = "datafusion_registration";
const SNAPSHOT_LOAD_FAILED: &str = "snapshot could not be loaded";
const EMPTY_TABLE_URI: &str = "table location must not be empty";
const INVALID_TABLE_URI: &str = "table location could not be parsed or normalized";
const ENGINE_CONSTRUCTION_FAILED: &str = "object store engine could not be constructed";
const S3_IMPLICIT_CREDENTIAL_HINT: &str = "S3 credential hint: no explicit S3 credentials were supplied through storage_options; local shells may need explicit AWS_ACCESS_KEY_ID, AWS_SECRET_ACCESS_KEY, optional AWS_SESSION_TOKEN, and AWS_REGION.";
#[derive(Clone, PartialEq, Eq)]
pub struct RegisteredSessionSource {
table: LazyTable,
source_uri: String,
snapshot_version: u64,
schema: SchemaRef,
protocol: DeltaProtocolReport,
phase_timings: Vec<PhaseTimingReport>,
}
impl RegisteredSessionSource {
pub(super) fn from_registered(
table: LazyTable,
registered: RegisteredDeltaSource,
phase_timings: Vec<PhaseTimingReport>,
) -> Self {
Self {
table,
source_uri: registered.table_uri,
snapshot_version: registered.snapshot_version,
schema: registered.schema,
protocol: registered.protocol,
phase_timings,
}
}
#[must_use]
pub const fn table(&self) -> &LazyTable {
&self.table
}
#[must_use]
pub fn name(&self) -> &str {
self.table.name()
}
#[must_use]
pub fn source_uri(&self) -> &str {
&self.source_uri
}
#[must_use]
pub const fn snapshot_version(&self) -> u64 {
self.snapshot_version
}
#[must_use]
pub fn schema(&self) -> &SchemaRef {
&self.schema
}
#[must_use]
pub const fn protocol(&self) -> &DeltaProtocolReport {
&self.protocol
}
#[must_use]
pub fn phase_timings(&self) -> &[PhaseTimingReport] {
&self.phase_timings
}
}
impl fmt::Debug for RegisteredSessionSource {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("RegisteredSessionSource")
.field("table", &self.table)
.field("source_uri", &self.source_uri)
.field("snapshot_version", &self.snapshot_version)
.field("schema", &self.schema)
.field("protocol", &self.protocol)
.field("phase_timings", &self.phase_timings)
.finish()
}
}
impl DeltaFunnelSession {
pub async fn delta_lake(
&mut self,
source: DeltaSourceConfig,
) -> Result<LazyTable, DeltaFunnelError> {
self.validate_delta_source_registration(&source)?;
self.register_delta_source(source, None).await
}
pub(crate) async fn delta_lake_with_progress(
&mut self,
source: DeltaSourceConfig,
reporter: ProgressReporter,
) -> Result<LazyTable, DeltaFunnelError> {
self.validate_delta_source_registration(&source)?;
reporter.emit(&ProgressEvent::started(
ProgressOperation::RegisterDeltaSource,
));
let result = self.register_delta_source(source, Some(&reporter)).await;
reporter.emit(&if result.is_ok() {
ProgressEvent::completed()
} else {
ProgressEvent::failed()
});
result
}
async fn register_delta_source(
&mut self,
source: DeltaSourceConfig,
reporter: Option<&ProgressReporter>,
) -> Result<LazyTable, DeltaFunnelError> {
let mut phase_timings = Vec::new();
emit_registration_phase(reporter, ProgressPhase::LoadingDeltaMetadata);
let source_timer = PhaseTimer::start(SOURCE_LOADING_PHASE);
let (source_name, snapshot) =
match load_delta_snapshot(source, self.options.provider_scan_options()).await {
Ok(loaded) => {
phase_timings.push(source_timer.completed());
loaded
}
Err(error) => {
phase_timings.push(source_timer.failed());
return Err(error);
}
};
emit_registration_phase(reporter, ProgressPhase::ValidatingDeltaProtocol);
let preflight_timer = PhaseTimer::start(PROTOCOL_PREFLIGHT_PHASE);
match validate_delta_protocol(&source_name, &snapshot) {
Ok(()) => {
phase_timings.push(preflight_timer.completed());
}
Err(error) => {
phase_timings.push(preflight_timer.failed());
return Err(error);
}
}
let registration_timer = PhaseTimer::start(DATAFUSION_REGISTRATION_PHASE);
emit_registration_phase(reporter, ProgressPhase::PreparingDeltaProvider);
let table_uri = snapshot.table_url().to_owned();
let table = match snapshot.into_table() {
Ok(table) => table,
Err(error) => {
phase_timings.push(registration_timer.failed());
return Err(map_load_error(&source_name, &table_uri, None, error));
}
};
let registered = match register_delta_source_with_scan_options(
&self.context,
source_name,
table,
ScanOptions {
execution_options: self.options.provider_scan_options(),
target_partitions: None,
intra_file_repartitioning: self.options.provider_file_repartitioning(),
use_arrow_view_types: self.options.provider_use_view_types(),
},
reporter,
) {
Ok(registered) => {
phase_timings.push(registration_timer.completed());
registered
}
Err(error) => {
phase_timings.push(registration_timer.failed());
return Err(error);
}
};
let table = self.allocate_delta_source_table(registered.name.clone());
let session_source =
RegisteredSessionSource::from_registered(table.clone(), registered, phase_timings);
self.sources.push(session_source);
Ok(table)
}
fn validate_delta_source_registration(
&self,
source: &DeltaSourceConfig,
) -> Result<(), DeltaFunnelError> {
validate_table_source_names([source.name.as_str()])?;
self.reject_registered_alias_name(&source.name)?;
reject_existing_delta_registration_name(&self.context, &source.name, &source.table_uri)
}
fn allocate_delta_source_table(&mut self, name: String) -> LazyTable {
let id = self.next_table_id;
self.next_table_id = self.next_table_id.saturating_add(1);
LazyTable::delta_source(id, name)
}
}
async fn load_delta_snapshot(
mut source: DeltaSourceConfig,
execution_options: delta_arrow_reader::DeltaScanExecutionOptions,
) -> Result<(String, DeltaTableSnapshot), DeltaFunnelError> {
source.apply_environment_storage_options();
let source_name = source.name.clone();
let table_uri = source.table_uri.clone();
let auth_mode = source.s3_auth_mode_hint();
observability::source_loading_started(&source_name, auth_mode.map(S3AuthModeHint::as_str));
let selection = source.version.map_or(
DeltaSnapshotSelection::Latest,
DeltaSnapshotSelection::Version,
);
let result = DeltaTableBuilder::new(&source.table_uri)
.with_storage_options(source.storage_options)
.with_snapshot_selection(selection)
.with_execution_options(execution_options)
.load_snapshot()
.await
.map_err(|error| map_load_error(&source_name, &table_uri, auth_mode, error));
match &result {
Ok(snapshot) => observability::source_loading_completed(
&source_name,
snapshot.version(),
auth_mode.map(S3AuthModeHint::as_str),
),
Err(error) => observability::source_loading_failed(
&source_name,
error,
auth_mode.map(S3AuthModeHint::as_str),
),
}
result.map(|snapshot| (source_name, snapshot))
}
fn validate_delta_protocol(
source_name: &str,
snapshot: &DeltaTableSnapshot,
) -> Result<(), DeltaFunnelError> {
observability::protocol_preflight_started(source_name, snapshot.version());
let result = validate_delta_table_snapshot_protocol(source_name, snapshot);
match &result {
Ok(()) => observability::protocol_preflight_completed(source_name, snapshot.version()),
Err(error) => {
observability::protocol_preflight_failed(source_name, snapshot.version(), error);
}
}
result
}
fn map_load_error(
source_name: &str,
table_uri: &str,
auth_mode: Option<S3AuthModeHint>,
error: DeltaReaderError,
) -> DeltaFunnelError {
match error.phase() {
DeltaReaderPhase::Configuration => DeltaFunnelError::Config {
message: error.to_string(),
},
DeltaReaderPhase::TableLocation => DeltaFunnelError::InvalidSourceUri {
reason: if table_uri.trim().is_empty() {
EMPTY_TABLE_URI
} else {
INVALID_TABLE_URI
},
},
DeltaReaderPhase::Storage => DeltaFunnelError::DeltaSourceEngine {
reason: ENGINE_CONSTRUCTION_FAILED,
},
DeltaReaderPhase::Snapshot => DeltaFunnelError::DeltaSnapshotLoad {
reason: snapshot_load_failed_reason(&error, auth_mode),
},
DeltaReaderPhase::Schema => DeltaFunnelError::DeltaSourceSchema {
source_name: source_name.to_owned(),
table_uri: table_uri.to_owned(),
reason: reader_error_source_reason(&error),
},
_ => DeltaFunnelError::DeltaSnapshotLoad {
reason: error.to_string(),
},
}
}
fn snapshot_load_failed_reason(
error: &DeltaReaderError,
auth_mode: Option<S3AuthModeHint>,
) -> String {
let cause = reader_error_source_reason(error);
let mut reason = format!(
"{SNAPSHOT_LOAD_FAILED}: {}",
sanitize_snapshot_load_cause(&cause)
);
if auth_mode == Some(S3AuthModeHint::ImplicitProviderChain) {
reason.push(' ');
reason.push_str(S3_IMPLICIT_CREDENTIAL_HINT);
}
reason
}
fn reader_error_source_reason(error: &DeltaReaderError) -> String {
error
.source()
.map_or_else(|| error.to_string(), ToString::to_string)
}
fn sanitize_snapshot_load_cause(cause: &str) -> String {
cause
.split_whitespace()
.map(|token| {
if token.contains("://") {
sanitize_uri_for_display(token)
} else {
sanitize_text_for_display(token)
}
})
.collect::<Vec<_>>()
.join(" ")
}
fn emit_registration_phase(reporter: Option<&ProgressReporter>, phase: ProgressPhase) {
if let Some(reporter) = reporter {
reporter.emit(&ProgressEvent::phase_changed(phase, None));
}
}
#[cfg(test)]
mod tests {
use std::sync::{Arc, Mutex};
use datafusion::{
arrow::datatypes::{DataType, Schema},
datasource::empty::EmptyTable,
};
use super::{
DATAFUSION_REGISTRATION_PHASE, EMPTY_TABLE_URI, PROTOCOL_PREFLIGHT_PHASE,
S3_IMPLICIT_CREDENTIAL_HINT, SOURCE_LOADING_PHASE, snapshot_load_failed_reason,
};
use crate::{
DeltaFunnelError, DeltaSourceConfig, QueryOptions,
progress::{
ProgressEvent, ProgressEventKind, ProgressOperation, ProgressPhase, ProgressReporter,
},
query_engine::datafusion::test_support::{
FailsOnCustomersSchemaProvider, INVALID_NESTED_IDS_SCHEMA_FIELDS_JSON,
SingleSchemaCatalogProvider,
},
};
use delta_arrow_reader::{
DeltaScanExecutionOptions, DeltaStorageOptions, ParquetReaderBackend,
};
use super::super::super::{
DeltaFunnelSession, LazyTableKind, SessionOptions, SourceUsageStatus,
test_support::DeltaLogTable,
};
const UNSUPPORTED_PROTOCOL_JSON: &str =
r#"{"protocol":{"minReaderVersion":99,"minWriterVersion":2}}"#;
const UNSUPPORTED_FEATURE_PROTOCOL_JSON: &str = r#"{"protocol":{"minReaderVersion":3,"minWriterVersion":7,"readerFeatures":["madeUpFeature","laterFeature"],"writerFeatures":["madeUpFeature","laterFeature"]}}"#;
#[tokio::test]
async fn snapshot_error_keeps_the_s3_hint_only_for_implicit_auth() {
let error = delta_arrow_reader::DeltaTableBuilder::new("")
.load_table()
.await
.expect_err("blank URI should fail");
let implicit = snapshot_load_failed_reason(
&error,
Some(crate::table_formats::S3AuthModeHint::ImplicitProviderChain),
);
let explicit = snapshot_load_failed_reason(
&error,
Some(crate::table_formats::S3AuthModeHint::ExplicitStatic),
);
let non_s3 = snapshot_load_failed_reason(&error, None);
assert!(implicit.contains(S3_IMPLICIT_CREDENTIAL_HINT));
assert!(!explicit.contains(S3_IMPLICIT_CREDENTIAL_HINT));
assert!(!non_s3.contains(S3_IMPLICIT_CREDENTIAL_HINT));
}
fn recording_reporter() -> (ProgressReporter, Arc<Mutex<Vec<ProgressEvent>>>) {
let events = Arc::new(Mutex::new(Vec::new()));
let recorded = Arc::clone(&events);
let reporter = ProgressReporter::new(move |event| {
if let Ok(mut events) = recorded.lock() {
events.push(event.clone());
}
});
(reporter, events)
}
async fn assert_source_loading_progress_failure(
source: DeltaSourceConfig,
error_matches: impl FnOnce(&DeltaFunnelError) -> bool,
) -> Result<(), Box<dyn std::error::Error>> {
let source_name = source.name.clone();
let mut session = DeltaFunnelSession::new(SessionOptions::default())?;
let (reporter, events) = recording_reporter();
let result = session.delta_lake_with_progress(source, reporter).await;
let error = result
.as_ref()
.err()
.ok_or("expected source loading error")?;
assert!(error_matches(error));
let events = events.lock().map_err(|_| "progress events lock poisoned")?;
assert_eq!(events.len(), 3);
assert_eq!(events[0].kind(), ProgressEventKind::Started);
assert_eq!(events[1].phase(), Some(ProgressPhase::LoadingDeltaMetadata));
assert_eq!(events[2].kind(), ProgressEventKind::Failed);
assert!(session.sources().is_empty());
assert_eq!(session.next_table_id(), 0);
assert!(!session.context().table_exist(&source_name)?);
Ok(())
}
#[tokio::test]
async fn delta_lake_progress_reports_ordered_registration_lifecycle()
-> Result<(), Box<dyn std::error::Error>> {
let table = DeltaLogTable::new("orders-progress")?;
let mut session = DeltaFunnelSession::new(SessionOptions::default())?;
let (reporter, events) = recording_reporter();
session
.delta_lake_with_progress(DeltaSourceConfig::new("orders", table.uri()), reporter)
.await?;
let events = events.lock().map_err(|_| "progress events lock poisoned")?;
assert_eq!(events.len(), 6);
assert_eq!(events[0].kind(), ProgressEventKind::Started);
assert_eq!(
events[0].operation(),
Some(ProgressOperation::RegisterDeltaSource)
);
assert_eq!(
events[1..5]
.iter()
.map(ProgressEvent::phase)
.collect::<Vec<_>>(),
vec![
Some(ProgressPhase::LoadingDeltaMetadata),
Some(ProgressPhase::ValidatingDeltaProtocol),
Some(ProgressPhase::PreparingDeltaProvider),
Some(ProgressPhase::RegisteringDeltaSource),
]
);
assert_eq!(events[5].kind(), ProgressEventKind::Completed);
assert!(events.iter().all(|event| event.output_name().is_none()));
assert!(events.iter().all(|event| event.files_total().is_none()));
assert!(events.iter().all(|event| event.rows().is_none()));
Ok(())
}
#[tokio::test]
async fn delta_lake_progress_preserves_each_reader_backend_registration()
-> Result<(), Box<dyn std::error::Error>> {
let table = DeltaLogTable::new("progress-reader-backends")?;
for backend in [
ParquetReaderBackend::DeltaKernel,
ParquetReaderBackend::Direct,
] {
let provider_options = DeltaScanExecutionOptions::new().with_parquet_backend(backend);
let mut ordinary = DeltaFunnelSession::new(
SessionOptions::new().with_provider_scan_options(provider_options),
)?;
let mut reported = DeltaFunnelSession::new(
SessionOptions::new().with_provider_scan_options(provider_options),
)?;
let source_uri = table.uri();
let ordinary_table = ordinary
.delta_lake(DeltaSourceConfig::new("orders", source_uri.clone()))
.await?;
let (reporter, events) = recording_reporter();
let reported_table = reported
.delta_lake_with_progress(DeltaSourceConfig::new("orders", source_uri), reporter)
.await?;
assert_eq!(reported_table, ordinary_table);
let ordinary_source = ordinary
.registered_source("orders")
.ok_or("ordinary source was not registered")?;
let reported_source = reported
.registered_source("orders")
.ok_or("reported source was not registered")?;
assert_eq!(reported_source.source_uri(), ordinary_source.source_uri());
assert_eq!(
reported_source.snapshot_version(),
ordinary_source.snapshot_version()
);
assert_eq!(reported_source.schema(), ordinary_source.schema());
assert_eq!(reported_source.protocol(), ordinary_source.protocol());
assert_eq!(
reported.source_reports()[0].scheduling().reader_backend(),
backend
);
let events = events.lock().map_err(|_| "progress events lock poisoned")?;
assert_eq!(
events
.iter()
.filter_map(ProgressEvent::phase)
.collect::<Vec<_>>(),
[
ProgressPhase::LoadingDeltaMetadata,
ProgressPhase::ValidatingDeltaProtocol,
ProgressPhase::PreparingDeltaProvider,
ProgressPhase::RegisteringDeltaSource,
]
);
assert_eq!(
events.last().map(ProgressEvent::kind),
Some(ProgressEventKind::Completed)
);
}
Ok(())
}
#[tokio::test]
async fn delta_provider_view_policy_defaults_to_standard_and_can_enable_views()
-> Result<(), Box<dyn std::error::Error>> {
let table = DeltaLogTable::new("provider-view-policy")?;
for (use_view_types, expected) in [(false, DataType::Utf8), (true, DataType::Utf8View)] {
let mut session = DeltaFunnelSession::new(
SessionOptions::new().with_provider_use_view_types(use_view_types),
)?;
session
.delta_lake(DeltaSourceConfig::new("orders", table.uri()))
.await?;
let provider = session.context().table_provider("orders").await?;
assert_eq!(
provider
.schema()
.field_with_name("customer_name")?
.data_type(),
&expected
);
}
Ok(())
}
#[tokio::test]
async fn delta_lake_progress_covers_source_loading_failure_variants()
-> Result<(), Box<dyn std::error::Error>> {
assert_source_loading_progress_failure(
DeltaSourceConfig::new("orders", ""),
|error| matches!(error, DeltaFunnelError::InvalidSourceUri { reason } if *reason == EMPTY_TABLE_URI),
)
.await?;
assert_source_loading_progress_failure(
DeltaSourceConfig::new("orders", "ftp://example.com/table"),
|error| matches!(error, DeltaFunnelError::DeltaSourceEngine { .. }),
)
.await?;
let table = DeltaLogTable::new("progress-missing-snapshot")?;
assert_source_loading_progress_failure(
DeltaSourceConfig::new("orders", table.uri()).with_version(Some(999)),
|error| matches!(error, DeltaFunnelError::DeltaSnapshotLoad { .. }),
)
.await?;
Ok(())
}
#[tokio::test]
async fn delta_lake_progress_reports_protocol_failure_after_validation_starts()
-> Result<(), Box<dyn std::error::Error>> {
let table =
DeltaLogTable::new_with_protocol("progress-protocol", UNSUPPORTED_PROTOCOL_JSON)?;
let mut session = DeltaFunnelSession::new(SessionOptions::default())?;
let (reporter, events) = recording_reporter();
let result = session
.delta_lake_with_progress(DeltaSourceConfig::new("orders", table.uri()), reporter)
.await;
assert!(matches!(
result,
Err(DeltaFunnelError::DeltaProtocolCompatibility { .. })
));
let events = events.lock().map_err(|_| "progress events lock poisoned")?;
assert_eq!(events.len(), 4);
assert_eq!(events[0].kind(), ProgressEventKind::Started);
assert_eq!(events[1].phase(), Some(ProgressPhase::LoadingDeltaMetadata));
assert_eq!(
events[2].phase(),
Some(ProgressPhase::ValidatingDeltaProtocol)
);
assert_eq!(events[3].kind(), ProgressEventKind::Failed);
Ok(())
}
#[tokio::test]
async fn delta_lake_progress_reports_provider_preparation_failure()
-> Result<(), Box<dyn std::error::Error>> {
let table = DeltaLogTable::new_with_schema(
"progress-provider",
INVALID_NESTED_IDS_SCHEMA_FIELDS_JSON,
)?;
let mut session = DeltaFunnelSession::new(SessionOptions::default())?;
let (reporter, events) = recording_reporter();
let result = session
.delta_lake_with_progress(DeltaSourceConfig::new("orders", table.uri()), reporter)
.await;
assert!(matches!(
result,
Err(DeltaFunnelError::DeltaSourceSchema { reason, .. })
if reason.contains("bad_array")
&& reason.contains("delta.columnMapping.nested.ids")
));
let events = events.lock().map_err(|_| "progress events lock poisoned")?;
assert_eq!(events.len(), 5);
assert_eq!(
events[3].phase(),
Some(ProgressPhase::PreparingDeltaProvider)
);
assert_eq!(events[4].kind(), ProgressEventKind::Failed);
assert!(
events
.iter()
.all(|event| event.phase() != Some(ProgressPhase::RegisteringDeltaSource))
);
Ok(())
}
#[tokio::test]
async fn delta_lake_progress_reports_catalog_registration_failure()
-> Result<(), Box<dyn std::error::Error>> {
let source_table = DeltaLogTable::new("progress-catalog")?;
let source_uri = source_table.uri();
let mut session = DeltaFunnelSession::new(SessionOptions::default())?;
let failing_schema = Arc::new(FailsOnCustomersSchemaProvider::default());
let schema: Arc<dyn datafusion::catalog::SchemaProvider> = failing_schema.clone();
session.context().register_catalog(
"datafusion",
Arc::new(SingleSchemaCatalogProvider::new(schema)),
);
let (reporter, events) = recording_reporter();
let result = session
.delta_lake_with_progress(
DeltaSourceConfig::new("customers", source_uri.clone()),
reporter,
)
.await;
assert!(matches!(
result,
Err(DeltaFunnelError::DataFusionRegistration { .. })
));
{
let events = events.lock().map_err(|_| "progress events lock poisoned")?;
assert_eq!(events.len(), 6);
assert_eq!(
events[4].phase(),
Some(ProgressPhase::RegisteringDeltaSource)
);
assert_eq!(events[5].kind(), ProgressEventKind::Failed);
}
assert!(session.sources().is_empty());
assert_eq!(session.next_table_id(), 0);
assert!(!session.context().table_exist("customers")?);
failing_schema.allow_customers();
let (reporter, retry_events) = recording_reporter();
let table = session
.delta_lake_with_progress(DeltaSourceConfig::new("customers", source_uri), reporter)
.await?;
assert_eq!(table.id(), 0);
assert_eq!(session.sources().len(), 1);
assert!(session.context().table_exist("customers")?);
let retry_events = retry_events
.lock()
.map_err(|_| "progress events lock poisoned")?;
assert_eq!(retry_events.len(), 6);
assert_eq!(retry_events[0].kind(), ProgressEventKind::Started);
assert_eq!(retry_events[5].kind(), ProgressEventKind::Completed);
Ok(())
}
#[tokio::test]
async fn delta_lake_progress_does_not_start_for_local_registration_errors()
-> Result<(), Box<dyn std::error::Error>> {
let table = DeltaLogTable::new("orders-conflict")?;
let mut session = DeltaFunnelSession::new(SessionOptions::default())?;
session
.delta_lake(DeltaSourceConfig::new("orders", table.uri()))
.await?;
for source in [
DeltaSourceConfig::new("select", ""),
DeltaSourceConfig::new("ORDERS", ""),
] {
let (reporter, events) = recording_reporter();
assert!(
session
.delta_lake_with_progress(source, reporter)
.await
.is_err()
);
assert!(
events
.lock()
.map_err(|_| "progress events lock poisoned")?
.is_empty()
);
}
Ok(())
}
#[tokio::test]
async fn delta_lake_progress_does_not_start_for_datafusion_catalog_conflict()
-> Result<(), Box<dyn std::error::Error>> {
let mut session = DeltaFunnelSession::new(SessionOptions::default())?;
session.context().register_table(
"ExistingOrders",
Arc::new(EmptyTable::new(Arc::new(Schema::empty()))),
)?;
let (reporter, events) = recording_reporter();
let result = session
.delta_lake_with_progress(
DeltaSourceConfig::new("existingorders", "secret://must-not-load"),
reporter,
)
.await;
assert!(matches!(
result,
Err(DeltaFunnelError::DataFusionRegistration { .. })
));
assert!(
events
.lock()
.map_err(|_| "progress events lock poisoned")?
.is_empty()
);
Ok(())
}
#[tokio::test]
async fn delta_lake_registers_source_and_returns_lazy_table()
-> Result<(), Box<dyn std::error::Error>> {
let table = DeltaLogTable::new("orders")?;
let mut session = DeltaFunnelSession::new(SessionOptions::default())?;
let lazy = session
.delta_lake(DeltaSourceConfig::new("orders", table.uri()))
.await?;
assert_eq!(lazy.id(), 0);
assert_eq!(lazy.kind(), LazyTableKind::DeltaSource);
assert_eq!(lazy.name(), "orders");
assert_eq!(session.next_table_id(), 1);
assert_eq!(session.sources().len(), 1);
let registered = session
.registered_source("ORDERS")
.ok_or("expected registered source")?;
assert_eq!(registered.table(), &lazy);
assert_eq!(registered.name(), "orders");
assert!(registered.source_uri().starts_with("file://"));
assert_eq!(registered.snapshot_version(), 1);
assert_eq!(registered.protocol().source_name, "orders");
assert_eq!(registered.schema().fields().len(), 2);
let source_reports = session.source_reports();
assert_eq!(source_reports.len(), 1);
let report = &source_reports[0];
assert_eq!(report.source_name(), "orders");
assert_eq!(report.source_uri(), registered.source_uri());
assert_eq!(report.snapshot_version(), 1);
assert_eq!(report.protocol().source_name, "orders");
assert_eq!(report.scheduling().query_target_partitions(), None);
assert_eq!(
report.scheduling().reader_backend(),
ParquetReaderBackend::Direct
);
assert_eq!(
report.scheduling().max_concurrent_file_reads_per_scan(),
None
);
assert_eq!(report.file_count(), crate::FileCount::unavailable());
assert_eq!(
report.file_count_reason(),
Some(crate::ReportReasonCode::CostAvoidance)
);
assert!(!report.scan_metadata_exhausted());
assert_eq!(report.usage_status(), SourceUsageStatus::Unknown);
assert!(report.used_by_output_names().is_empty());
assert!(report.provider_read_stats().is_none());
assert_eq!(
report.provider_stats_reason(),
Some(crate::ReportReasonCode::NotExecuted)
);
let phase_timings = report.phase_timings();
assert_eq!(
phase_timings
.iter()
.map(crate::PhaseTimingReport::phase_name)
.collect::<Vec<_>>(),
vec![
SOURCE_LOADING_PHASE,
PROTOCOL_PREFLIGHT_PHASE,
DATAFUSION_REGISTRATION_PHASE
]
);
assert!(
phase_timings
.iter()
.all(|timing| timing.status().is_completed())
);
assert!(
phase_timings
.iter()
.all(|timing| timing.elapsed_micros().is_some())
);
Ok(())
}
#[tokio::test]
async fn delta_lake_registers_multiple_distinct_sources()
-> Result<(), Box<dyn std::error::Error>> {
let orders = DeltaLogTable::new("orders")?;
let customers = DeltaLogTable::new("customers")?;
let mut session = DeltaFunnelSession::new(SessionOptions::default())?;
let orders = session
.delta_lake(DeltaSourceConfig::new("orders", orders.uri()))
.await?;
let customers = session
.delta_lake(DeltaSourceConfig::new("customers", customers.uri()))
.await?;
assert_eq!(orders.id(), 0);
assert_eq!(customers.id(), 1);
assert_eq!(session.sources().len(), 2);
assert!(session.registered_source("orders").is_some());
assert!(session.registered_source("customers").is_some());
Ok(())
}
#[tokio::test]
async fn duplicate_source_alias_fails_before_loading_second_source()
-> Result<(), Box<dyn std::error::Error>> {
let table = DeltaLogTable::new("orders")?;
let mut session = DeltaFunnelSession::new(SessionOptions::default())?;
session
.delta_lake(DeltaSourceConfig::new("orders", table.uri()))
.await?;
let error = session
.delta_lake(DeltaSourceConfig::new("ORDERS", ""))
.await;
assert!(matches!(
error,
Err(DeltaFunnelError::DuplicateSourceName { name }) if name == "ORDERS"
));
assert_eq!(session.sources().len(), 1);
assert_eq!(session.next_table_id(), 1);
Ok(())
}
#[tokio::test]
async fn invalid_source_alias_fails_before_registration() -> Result<(), DeltaFunnelError> {
let mut session = DeltaFunnelSession::new(SessionOptions::default())?;
let error = session
.delta_lake(DeltaSourceConfig::new("select", ""))
.await;
assert!(matches!(
error,
Err(DeltaFunnelError::InvalidSourceName { name, .. }) if name == "select"
));
assert!(session.sources().is_empty());
assert_eq!(session.next_table_id(), 0);
Ok(())
}
#[tokio::test]
async fn protocol_preflight_failure_does_not_register_source()
-> Result<(), Box<dyn std::error::Error>> {
let table = DeltaLogTable::new_with_protocol("unsupported", UNSUPPORTED_PROTOCOL_JSON)?;
let mut session = DeltaFunnelSession::new(SessionOptions::default())?;
let error = session
.delta_lake(DeltaSourceConfig::new("unsupported", table.uri()))
.await;
let display = format!("{}", error.as_ref().err().ok_or("expected error")?);
assert!(display.contains("unsupported"));
assert!(display.contains("unsupported Delta minReaderVersion"));
assert!(matches!(
error,
Err(DeltaFunnelError::DeltaProtocolCompatibility { .. })
));
assert!(session.sources().is_empty());
assert_eq!(session.next_table_id(), 0);
Ok(())
}
#[tokio::test]
async fn protocol_preflight_reports_the_first_unsupported_reader_feature()
-> Result<(), Box<dyn std::error::Error>> {
let table = DeltaLogTable::new_with_protocol(
"unsupported-feature",
UNSUPPORTED_FEATURE_PROTOCOL_JSON,
)?;
let mut session = DeltaFunnelSession::new(SessionOptions::default())?;
let error = session
.delta_lake(DeltaSourceConfig::new("unsupported", table.uri()))
.await
.err()
.ok_or("expected protocol compatibility error")?;
let DeltaFunnelError::DeltaProtocolCompatibility { reason, .. } = error else {
return Err("expected protocol compatibility error".into());
};
assert_eq!(reason, "unsupported Delta reader feature `madeUpFeature`");
assert!(session.sources().is_empty());
Ok(())
}
#[tokio::test]
async fn protocol_preflight_failure_redacts_secret_uri_parts()
-> Result<(), Box<dyn std::error::Error>> {
let table = DeltaLogTable::new_with_protocol("unsupported", UNSUPPORTED_PROTOCOL_JSON)?;
let mut session = DeltaFunnelSession::new(SessionOptions::default())?;
let error = session
.delta_lake(DeltaSourceConfig::new(
"unsupported",
table.file_uri_with_secret_parts()?,
))
.await
.map(|_| ())
.map_err(|error| error.to_string());
assert!(
matches!(error, Err(display) if display.contains("unsupported")
&& display.contains("unsupported Delta minReaderVersion")
&& !display.contains("super-secret")
&& !display.contains("debug-secret")
&& !display.contains("token"))
);
assert!(session.sources().is_empty());
assert_eq!(session.next_table_id(), 0);
Ok(())
}
#[tokio::test]
async fn protocol_preflight_failure_does_not_leak_datafusion_alias()
-> Result<(), Box<dyn std::error::Error>> {
let table = DeltaLogTable::new_with_protocol("unsupported", UNSUPPORTED_PROTOCOL_JSON)?;
let mut session = DeltaFunnelSession::new(SessionOptions::default())?;
let error = session
.delta_lake(DeltaSourceConfig::new("unsupported", table.uri()))
.await;
assert!(matches!(
error,
Err(DeltaFunnelError::DeltaProtocolCompatibility { .. })
));
assert!(session.context().table("unsupported").await.is_err());
assert!(session.sources().is_empty());
assert_eq!(session.next_table_id(), 0);
Ok(())
}
#[tokio::test]
async fn registered_source_sql_analysis_does_not_read_data_files_for_rows()
-> Result<(), Box<dyn std::error::Error>> {
let table = DeltaLogTable::new("orders")?;
let mut session = DeltaFunnelSession::new(SessionOptions::default())?;
session
.delta_lake(DeltaSourceConfig::new("orders", table.uri()))
.await?;
let dataframe = session
.context()
.sql("select id, customer_name from orders")
.await?;
let schema = dataframe.schema();
assert_eq!(schema.fields().len(), 2);
assert_eq!(schema.field(0).name(), "id");
assert_eq!(schema.field(1).name(), "customer_name");
assert_eq!(session.sources().len(), 1);
Ok(())
}
#[tokio::test]
async fn source_debug_does_not_expose_storage_option_values()
-> Result<(), Box<dyn std::error::Error>> {
let table = DeltaLogTable::new("storage-options")?;
let mut storage_options = DeltaStorageOptions::new();
storage_options.insert(
"AWS_SECRET_ACCESS_KEY".to_owned(),
"super-secret".to_owned(),
);
let mut session = DeltaFunnelSession::new(SessionOptions::default())?;
session
.delta_lake(
DeltaSourceConfig::new("orders", table.uri()).with_storage_options(storage_options),
)
.await?;
let debug = format!("{session:?}");
assert!(debug.contains("orders"));
assert!(!debug.contains("super-secret"));
assert!(!debug.contains("AWS_SECRET_ACCESS_KEY"));
let report_debug = format!("{:?}", session.source_reports());
assert!(report_debug.contains("orders"));
assert!(!report_debug.contains("super-secret"));
assert!(!report_debug.contains("AWS_SECRET_ACCESS_KEY"));
Ok(())
}
#[tokio::test]
async fn source_registration_honors_configured_provider_options()
-> Result<(), Box<dyn std::error::Error>> {
let table = DeltaLogTable::new("configured-provider")?;
let provider_scan_options = DeltaScanExecutionOptions::new()
.with_parquet_backend(ParquetReaderBackend::DeltaKernel)
.with_max_concurrent_file_reads_per_scan(Some(2))?
.with_max_concurrent_file_reads_per_partition(1)?
.with_output_buffer_batches_per_partition(3)?
.with_prefetch_files_per_partition(0)
.with_parquet_metadata_size_hint_bytes(Some(16_384))?;
let mut session = DeltaFunnelSession::new(
SessionOptions::new()
.with_query_options(QueryOptions {
target_partitions: Some(4),
output_batch_size: None,
})
.with_provider_scan_options(provider_scan_options),
)?;
session
.delta_lake(DeltaSourceConfig::new("orders", table.uri()))
.await?;
assert_eq!(session.sources().len(), 1);
assert!(session.registered_source("orders").is_some());
let reports = session.source_reports();
assert_eq!(reports.len(), 1);
let scheduling = reports[0].scheduling();
assert_eq!(scheduling.query_target_partitions(), Some(4));
assert_eq!(
scheduling.reader_backend(),
ParquetReaderBackend::DeltaKernel
);
assert_eq!(scheduling.max_concurrent_file_reads_per_scan(), Some(2));
assert_eq!(scheduling.max_concurrent_file_reads_per_partition(), 1);
assert_eq!(scheduling.output_buffer_capacity_per_partition(), 3);
assert_eq!(
scheduling.native_async_prefetch_file_count_per_partition(),
0
);
assert_eq!(scheduling.parquet_metadata_size_hint(), Some(16_384));
Ok(())
}
}