use std::sync::{Arc, LazyLock};
use tracing::{info, instrument};
use url::Url;
use super::LogSegment;
#[cfg(all(feature = "adaptive-metadata-in-dev", feature = "declarative-plans"))]
use crate::actions::CHECKPOINT_ACTION_NAME;
#[cfg(feature = "adaptive-metadata-in-dev")]
use crate::actions::{CheckpointAction, CHECKPOINT_ACTION_FIELD};
use crate::actions::{Metadata, Protocol, METADATA_FIELD, PROTOCOL_FIELD};
#[cfg(feature = "declarative-plans")]
use crate::actions::{METADATA_NAME, PROTOCOL_NAME};
use crate::crc::Crc;
use crate::engine_data::{GetData, RowVisitor, TypedGetData as _};
use crate::log_replay::ActionsBatch;
use crate::metrics::ProtocolMetadataSource;
use crate::path::ParsedLogPath;
#[cfg(feature = "declarative-plans")]
use crate::plans::ir::nodes::Agg;
#[cfg(feature = "declarative-plans")]
use crate::plans::ir::nodes::FileType;
#[cfg(feature = "declarative-plans")]
use crate::plans::{Operation, PlanBuilder, PlanExecutor};
use crate::schema::{
column_name, schema_ref, ColumnName, ColumnNamesAndTypes, DataType, MetadataColumnSpec,
StructField, StructType,
};
use crate::{DeltaResult, Engine, EngineData, Error, Version};
impl LogSegment {
pub(crate) fn read_protocol_metadata(
&self,
engine: &dyn Engine,
crc: Option<&Arc<Crc>>,
) -> DeltaResult<(Metadata, Protocol, ProtocolMetadataSource)> {
match self.read_protocol_metadata_opt(engine, crc)? {
(Some(m), Some(p), source) => Ok((m, p, source)),
(None, Some(_), _) => Err(Error::MissingMetadata),
(Some(_), None, _) => Err(Error::MissingProtocol),
(None, None, _) => Err(Error::MissingMetadataAndProtocol),
}
}
#[instrument(name = "log_seg.load_p_m", skip_all, err)]
pub(crate) fn read_protocol_metadata_opt(
&self,
engine: &dyn Engine,
crc: Option<&Arc<Crc>>,
) -> DeltaResult<(Option<Metadata>, Option<Protocol>, ProtocolMetadataSource)> {
if let Some(crc) = crc.filter(|c| c.version == self.end_version) {
info!("P&M from CRC at target version {}", self.end_version);
return Ok((
Some(crc.metadata.clone()),
Some(crc.protocol.clone()),
ProtocolMetadataSource::CrcAtTarget,
));
}
if let Some(crc) = crc.filter(|c| c.version < self.end_version) {
info!(
"Pruning log segment to commits after CRC version {}",
crc.version
);
let pruned = self.segment_after_version(crc.version);
let PmCandidate {
metadata: metadata_opt,
protocol: protocol_opt,
} = pruned.replay_for_pm(engine)?;
let metadata_opt = metadata_opt
.filter(|(v, _)| *v > crc.version as i64)
.map(|(_, m)| m);
let protocol_opt = protocol_opt
.filter(|(v, _)| *v > crc.version as i64)
.map(|(_, p)| p);
if metadata_opt.is_some() && protocol_opt.is_some() {
info!("Found P&M from pruned log replay");
return Ok((
metadata_opt,
protocol_opt,
ProtocolMetadataSource::CrcSeededPmOnlyReplay,
));
}
info!("P&M fallback to CRC (no P&M changes after CRC version)");
return Ok((
metadata_opt.or_else(|| Some(crc.metadata.clone())),
protocol_opt.or_else(|| Some(crc.protocol.clone())),
ProtocolMetadataSource::CrcSeededPmOnlyReplay,
));
}
let PmCandidate {
metadata: metadata_opt,
protocol: protocol_opt,
} = self.replay_for_pm(engine)?;
Ok((
metadata_opt.map(|(_, m)| m),
protocol_opt.map(|(_, p)| p),
ProtocolMetadataSource::FullReplay,
))
}
fn replay_for_pm(&self, engine: &dyn Engine) -> DeltaResult<PmCandidate> {
#[cfg(feature = "declarative-plans")]
if let Some(executor) = engine.plan_executor() {
return resolve_pm_batches(self.read_pm_batches_via_plan(executor.as_ref())?);
}
resolve_pm_batches(self.read_pm_batches(engine)?)
}
#[cfg(feature = "declarative-plans")]
fn read_pm_batches_via_plan(
&self,
executor: &dyn PlanExecutor,
) -> DeltaResult<impl Iterator<Item = DeltaResult<VersionedBatch>> + Send> {
#[cfg(feature = "adaptive-metadata-in-dev")]
let versioned_schema = schema_ref! {
(&PROTOCOL_FIELD),
(&METADATA_FIELD),
(&CHECKPOINT_ACTION_FIELD),
not_null "version": LONG,
};
#[cfg(not(feature = "adaptive-metadata-in-dev"))]
let versioned_schema = schema_ref! {
(&PROTOCOL_FIELD),
(&METADATA_FIELD),
not_null "version": LONG,
};
let commit_files = self.commit_cover_version_tagged_scan_files()?;
let commits = PlanBuilder::scan_json(commit_files, &["version"], versioned_schema.clone())?;
let checkpoint = self
.checkpoint_version_tagged_scan_files()?
.map(|(file_type, checkpoint_files)| {
let scan = match file_type {
FileType::Json => PlanBuilder::scan_json,
FileType::Parquet => PlanBuilder::scan_parquet,
};
scan(checkpoint_files, &["version"], versioned_schema.clone())
})
.transpose()?;
let plan = PlanBuilder::union_all(std::iter::once(commits).chain(checkpoint))?
.aggregate_ungrouped(|a| {
let protocol = || column_name!(PROTOCOL_NAME);
let metadata = || column_name!(METADATA_NAME);
let version = || column_name!("version");
let a = a
.max_non_null_by(protocol(), protocol(), version())
.max_non_null_by(metadata(), metadata(), version())
.aggregate_as(
Agg::max_non_null_by(version(), protocol(), version()),
"protocol_version",
)
.aggregate_as(
Agg::max_non_null_by(version(), metadata(), version()),
"metadata_version",
);
#[cfg(feature = "adaptive-metadata-in-dev")]
let a = a.max_non_null_by(
column_name!(CHECKPOINT_ACTION_NAME),
column_name!(CHECKPOINT_ACTION_NAME),
version(),
);
a
})?
.build()?;
let batches = executor
.execute_op(Operation::QueryPlan(plan))?
.into_data()?
.map(|batch| {
let batch = ActionsBatch::new(batch?, true);
let (protocol_version, metadata_version) =
pm_versions_from_plan_output(batch.actions.as_ref())?;
Ok(VersionedBatch {
protocol_version,
metadata_version,
batch,
})
});
Ok(batches)
}
fn read_pm_batches(
&self,
engine: &dyn Engine,
) -> DeltaResult<impl Iterator<Item = DeltaResult<VersionedBatch>> + Send> {
let (commit_schema, checkpoint_schema) = pm_replay_schemas();
let file_column =
StructField::create_metadata_column("_file", MetadataColumnSpec::FilePath);
let commit_schema = Arc::new(StructType::try_new(
commit_schema.fields().cloned().chain([file_column]),
)?);
let checkpoint_version = self.checkpoint_version.map(|v| v as i64);
let batches = self
.read_actions_with_projected_checkpoint_actions(
engine,
commit_schema,
checkpoint_schema,
None,
None,
None,
None,
)?
.actions;
Ok(batches.map(move |batch| {
let batch = batch?;
let version = if batch.is_log_batch {
batch_version(batch.actions.as_ref())? as i64
} else {
checkpoint_version
.ok_or_else(|| Error::internal_error("checkpoint batch without a version"))?
};
Ok(VersionedBatch {
protocol_version: Some(version),
metadata_version: Some(version),
batch,
})
}))
}
}
struct PmCandidate {
protocol: Option<(i64, Protocol)>,
metadata: Option<(i64, Metadata)>,
}
struct VersionedBatch {
protocol_version: Option<i64>,
metadata_version: Option<i64>,
batch: ActionsBatch,
}
fn resolve_pm_batches(
batches: impl Iterator<Item = DeltaResult<VersionedBatch>>,
) -> DeltaResult<PmCandidate> {
let mut metadata: Option<(i64, Metadata)> = None;
let mut protocol: Option<(i64, Protocol)> = None;
for batch in batches {
let VersionedBatch {
protocol_version,
metadata_version,
batch,
} = batch?;
let batch_version = protocol_version.max(metadata_version);
let candidate = pm_candidate(&batch, protocol_version, metadata_version)?;
metadata = newer(metadata, candidate.metadata);
protocol = newer(protocol, candidate.protocol);
if is_final(&protocol, batch_version) && is_final(&metadata, batch_version) {
break;
}
}
Ok(PmCandidate { protocol, metadata })
}
fn batch_version(data: &dyn EngineData) -> DeltaResult<Version> {
#[derive(Default)]
struct FilePathVisitor {
file: Option<String>,
}
impl RowVisitor for FilePathVisitor {
fn selected_column_names_and_types(&self) -> (&'static [ColumnName], &'static [DataType]) {
static NAMES_AND_TYPES: LazyLock<ColumnNamesAndTypes> =
LazyLock::new(|| (vec![column_name!("_file")], vec![DataType::STRING]).into());
NAMES_AND_TYPES.as_ref()
}
fn visit<'a>(
&mut self,
row_count: usize,
getters: &[&'a dyn GetData<'a>],
) -> DeltaResult<()> {
if self.file.is_none() && row_count > 0 {
self.file = getters[0].get_opt(0, "_file")?;
}
Ok(())
}
}
let mut visitor = FilePathVisitor::default();
visitor.visit_rows_of(data)?;
let file = visitor
.file
.ok_or_else(|| Error::internal_error("commit batch missing _file column"))?;
let url = Url::parse(&file)
.map_err(|e| Error::internal_error(format!("batch has invalid _file {file}: {e}")))?;
ParsedLogPath::try_from(url)?
.map(|path| path.version)
.ok_or_else(|| Error::internal_error(format!("batch from non-log file {file}")))
}
fn is_final<T>(winner: &Option<(i64, T)>, batch_version: Option<i64>) -> bool {
match (winner, batch_version) {
(Some((version, _)), Some(bv)) => *version >= bv,
_ => false,
}
}
fn newer<T>(a: Option<(i64, T)>, b: Option<(i64, T)>) -> Option<(i64, T)> {
match (a, b) {
(Some(a), Some(b)) => {
if b.0 >= a.0 {
Some(b)
} else {
Some(a)
}
}
(a, b) => a.or(b),
}
}
fn pm_replay_schemas() -> (Arc<StructType>, Arc<StructType>) {
let checkpoint_schema = schema_ref! {
(&PROTOCOL_FIELD),
(&METADATA_FIELD),
};
#[cfg(feature = "adaptive-metadata-in-dev")]
let commit_schema = schema_ref! {
(&PROTOCOL_FIELD),
(&METADATA_FIELD),
(&CHECKPOINT_ACTION_FIELD),
};
#[cfg(not(feature = "adaptive-metadata-in-dev"))]
let commit_schema = checkpoint_schema.clone();
(commit_schema, checkpoint_schema)
}
fn pm_candidate(
batch: &ActionsBatch,
protocol_version: Option<i64>,
metadata_version: Option<i64>,
) -> DeltaResult<PmCandidate> {
let actions = batch.actions.as_ref();
let protocol = protocol_version.zip(Protocol::try_new_from_data(actions)?);
let metadata = metadata_version.zip(Metadata::try_new_from_data(actions)?);
let (checkpoint_protocol, checkpoint_metadata) = match checkpoint_pm(batch)? {
Some((version, p, m)) => (Some((version, p)), Some((version, m))),
None => (None, None),
};
Ok(PmCandidate {
protocol: newer(protocol, checkpoint_protocol),
metadata: newer(metadata, checkpoint_metadata),
})
}
fn checkpoint_pm(batch: &ActionsBatch) -> DeltaResult<Option<(i64, Protocol, Metadata)>> {
#[cfg(feature = "adaptive-metadata-in-dev")]
{
if !batch.is_log_batch {
return Ok(None);
}
let checkpoint = CheckpointAction::try_new_from_data(batch.actions.as_ref())?;
Ok(checkpoint.map(|checkpoint| {
(
checkpoint.version(),
checkpoint.protocol().clone(),
checkpoint.metadata().clone(),
)
}))
}
#[cfg(not(feature = "adaptive-metadata-in-dev"))]
{
let _ = batch;
Ok(None)
}
}
#[cfg(feature = "declarative-plans")]
fn pm_versions_from_plan_output(
actions: &dyn EngineData,
) -> DeltaResult<(Option<i64>, Option<i64>)> {
#[derive(Default)]
struct PmVersionsVisitor {
protocol: Option<i64>,
metadata: Option<i64>,
}
impl RowVisitor for PmVersionsVisitor {
fn selected_column_names_and_types(&self) -> (&'static [ColumnName], &'static [DataType]) {
static NAMES_AND_TYPES: LazyLock<ColumnNamesAndTypes> = LazyLock::new(|| {
(
vec![
column_name!("protocol_version"),
column_name!("metadata_version"),
],
vec![DataType::LONG, DataType::LONG],
)
.into()
});
NAMES_AND_TYPES.as_ref()
}
fn visit<'a>(
&mut self,
row_count: usize,
getters: &[&'a dyn GetData<'a>],
) -> DeltaResult<()> {
if row_count > 0 {
self.protocol = getters[0].get_opt(0, "protocol_version")?;
self.metadata = getters[1].get_opt(0, "metadata_version")?;
}
Ok(())
}
}
let mut visitor = PmVersionsVisitor::default();
visitor.visit_rows_of(actions)?;
Ok((visitor.protocol, visitor.metadata))
}
#[cfg(test)]
mod tests {
use std::path::PathBuf;
#[cfg(feature = "declarative-plans")]
use std::sync::Arc;
use itertools::Itertools;
use test_log::test;
use crate::engine::sync::SyncEngine;
#[cfg(feature = "declarative-plans")]
use crate::engine::test_delegating::DelegatingEngine;
#[cfg(feature = "declarative-plans")]
use crate::plans::{Operation, PlanExecutor, PlanResult};
use crate::Snapshot;
#[cfg(feature = "declarative-plans")]
use crate::{DeltaResult, Error};
#[cfg(feature = "declarative-plans")]
struct FailingPlanExecutor;
#[cfg(feature = "declarative-plans")]
impl PlanExecutor for FailingPlanExecutor {
fn execute_op(&self, _op: Operation) -> DeltaResult<PlanResult> {
Err(Error::generic("plan executor deliberately failed"))
}
}
#[test]
fn test_replay_for_metadata() {
let path = std::fs::canonicalize(PathBuf::from("./tests/data/parquet_row_group_skipping/"));
let url = url::Url::from_directory_path(path.unwrap()).unwrap();
let engine = SyncEngine::new();
let snapshot = Snapshot::builder_for(url).build(&engine).unwrap();
let data: Vec<_> = snapshot
.log_segment()
.read_pm_batches(&engine)
.unwrap()
.try_collect()
.unwrap();
assert_eq!(data.len(), 4);
}
#[test]
fn test_snapshot_build_via_plan_over_parquet_checkpoint_with_entries_named_maps() {
let path =
std::fs::canonicalize(PathBuf::from("./tests/data/app-txn-checkpoint/")).unwrap();
let url = url::Url::from_directory_path(path).unwrap();
let engine = SyncEngine::new();
let snapshot = Snapshot::builder_for(url).build(&engine).unwrap();
assert_eq!(snapshot.version(), 1);
assert_eq!(snapshot.schema().fields().count(), 3);
}
#[test]
fn test_snapshot_build_via_plan_over_parquet_checkpoint_with_item_named_arrays() {
let path = std::fs::canonicalize(PathBuf::from("./tests/data/parsed-stats/")).unwrap();
let url = url::Url::from_directory_path(path).unwrap();
let engine = SyncEngine::new();
let snapshot = Snapshot::builder_for(url).build(&engine).unwrap();
assert_eq!(snapshot.version(), 5);
assert_eq!(snapshot.schema().fields().count(), 5);
}
#[cfg(feature = "declarative-plans")]
#[test]
fn test_snapshot_build_via_failing_plan_executor_surfaces_error_without_fallback() {
let path =
std::fs::canonicalize(PathBuf::from("./tests/data/app-txn-checkpoint/")).unwrap();
let url = url::Url::from_directory_path(path).unwrap();
let engine = DelegatingEngine::new(Arc::new(SyncEngine::new()))
.with_plan_executor(Arc::new(FailingPlanExecutor));
let result = Snapshot::builder_for(url).build(&engine);
assert!(
result.is_err(),
"plan failure must surface, not fall back to legacy replay"
);
}
}