use std::fmt;
use std::ops::Range;
use std::sync::Arc;
use arrow::datatypes::SchemaRef;
use async_trait::async_trait;
use bytes::Bytes;
use datafusion::catalog::{Session, TableProvider};
use datafusion::datasource::TableType;
use datafusion::datasource::listing::PartitionedFile;
use datafusion::datasource::physical_plan::{
FileGroup, FileScanConfigBuilder, ParquetFileReaderFactory, ParquetSource,
};
use datafusion::datasource::source::DataSourceExec;
use datafusion::error::{DataFusionError, Result as DataFusionResult};
use datafusion::execution::object_store::ObjectStoreUrl;
use datafusion::logical_expr::{Expr, TableProviderFilterPushDown};
use datafusion::physical_plan::ExecutionPlan;
use datafusion::physical_plan::metrics::ExecutionPlanMetricsSet;
use futures::FutureExt;
use futures::future::{BoxFuture, ready};
use parquet::arrow::arrow_reader::ArrowReaderOptions;
use parquet::arrow::async_reader::AsyncFileReader;
use parquet::arrow::parquet_to_arrow_schema;
use parquet::errors::ParquetError;
use parquet::file::metadata::{ParquetMetaData, ParquetMetaDataReader};
use polyc_projection_artifact::{FleetRealm, RetainedProjectionFile, VisibleRealm};
use polyc_state::immutable::ContentReference;
use polyc_state::projection::artifact::ObjectNamespace;
use polyc_state::revision::JournalSource;
use tokio_util::sync::CancellationToken;
use super::{CoreExecutionError, CoreOperationContext, CoreRealm, operation_refusal};
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
pub(super) struct VerifiedFileIdentity {
namespace: ObjectNamespace,
key: ContentReference,
generation: u64,
realm: CoreRealm,
}
#[derive(Clone)]
pub(super) enum VerifiedCoreFile {
Visible {
file: RetainedProjectionFile<VisibleRealm>,
source: JournalSource,
},
Fleet {
file: RetainedProjectionFile<FleetRealm>,
source: JournalSource,
},
}
impl fmt::Debug for VerifiedCoreFile {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("VerifiedCoreFile")
.field(
"realm",
match self {
Self::Visible { .. } => &"visible",
Self::Fleet { .. } => &"fleet",
},
)
.field("byte_len", &self.byte_len())
.finish_non_exhaustive()
}
}
impl VerifiedCoreFile {
pub(super) const fn byte_len(&self) -> u64 {
match self {
Self::Visible { file, .. } => file.byte_len(),
Self::Fleet { file, .. } => file.byte_len(),
}
}
const fn declared_rows(&self) -> u64 {
match self {
Self::Visible { file, .. } => file.descriptor().row_count(),
Self::Fleet { file, .. } => file.descriptor().row_count(),
}
}
const fn source(&self) -> &JournalSource {
match self {
Self::Visible { source, .. } | Self::Fleet { source, .. } => source,
}
}
pub(super) fn identity(&self) -> VerifiedFileIdentity {
let (object, realm) = match self {
Self::Visible { file, .. } => (file.object(), CoreRealm::Visible),
Self::Fleet { file, .. } => (file.object(), CoreRealm::Fleet),
};
VerifiedFileIdentity {
namespace: object.namespace().clone(),
key: object.key().clone(),
generation: object.generation(),
realm,
}
}
fn read_range(
&self,
offset: u64,
len: u64,
) -> Result<Bytes, polyc_projection_artifact::ArtifactReadError> {
match self {
Self::Visible { file, .. } => file.slice(offset, len),
Self::Fleet { file, .. } => file.slice(offset, len),
}
}
fn bytes(&self) -> &Bytes {
match self {
Self::Visible { file, .. } => file.bytes(),
Self::Fleet { file, .. } => file.bytes(),
}
}
const fn physical(&self) -> polyc_state::projection::artifact::PhysicalDecodeEnvelope {
match self {
Self::Visible { file, .. } => file.descriptor().physical(),
Self::Fleet { file, .. } => file.descriptor().physical(),
}
}
}
#[derive(Clone)]
struct ExactFileExtension {
file: VerifiedCoreFile,
operation: Arc<CoreOperationContext>,
cancellation: CancellationToken,
}
#[derive(Debug)]
struct ExactParquetReaderFactory;
impl ParquetFileReaderFactory for ExactParquetReaderFactory {
fn create_reader(
&self,
_partition_index: usize,
file: PartitionedFile,
_metadata_size_hint: Option<usize>,
_metrics: &ExecutionPlanMetricsSet,
) -> DataFusionResult<Box<dyn AsyncFileReader + Send>> {
let extension = file
.extension::<ExactFileExtension>()
.cloned()
.ok_or_else(|| DataFusionError::Execution("exact file token is absent".to_owned()))?;
Ok(Box::new(ExactParquetReader::from(extension)))
}
}
struct ExactParquetReader {
capability: ExactFileExtension,
}
impl From<ExactFileExtension> for ExactParquetReader {
fn from(value: ExactFileExtension) -> Self {
Self { capability: value }
}
}
impl ExactParquetReader {
fn read(&self, range: Range<u64>) -> Result<Bytes, ParquetError> {
if range.start > range.end || range.end > self.capability.file.byte_len() {
return Err(parquet_error(
"Parquet requested an out-of-bound exact range",
));
}
if self.capability.cancellation.is_cancelled() {
return Err(parquet_error("exact Parquet read was cancelled"));
}
self.capability
.operation
.check()
.map_err(operation_refusal)
.map_err(|error| ParquetError::External(Box::new(error)))?;
self.capability
.file
.read_range(range.start, range.end - range.start)
.map_err(|error| parquet_error(&error.to_string()))
}
}
impl AsyncFileReader for ExactParquetReader {
fn get_bytes(&mut self, range: Range<u64>) -> BoxFuture<'_, Result<Bytes, ParquetError>> {
ready(self.read(range)).boxed()
}
fn get_byte_ranges(
&mut self,
ranges: Vec<Range<u64>>,
) -> BoxFuture<'_, Result<Vec<Bytes>, ParquetError>> {
ready(ranges.into_iter().map(|range| self.read(range)).collect()).boxed()
}
fn get_metadata<'a>(
&'a mut self,
options: Option<&'a ArrowReaderOptions>,
) -> BoxFuture<'a, Result<Arc<ParquetMetaData>, ParquetError>> {
let file_len = self.capability.file.byte_len();
async move {
let metadata_options = options.map(|value| value.metadata_options().clone());
let metadata = ParquetMetaDataReader::new()
.with_metadata_options(metadata_options)
.load_and_finish(self, file_len)
.await?;
Ok(Arc::new(metadata))
}
.boxed()
}
}
fn parquet_error(reason: &str) -> ParquetError {
ParquetError::General(reason.to_owned())
}
pub(super) async fn validate_parquet_file(
file: VerifiedCoreFile,
expected_schema: SchemaRef,
operation: Arc<CoreOperationContext>,
cancellation: CancellationToken,
) -> Result<(), CoreExecutionError> {
let expected_rows = file.declared_rows();
let expected_source = file.source().clone();
polyc_projection_artifact::parquet_profile::inspect_exact(
file.bytes(),
expected_schema.as_ref(),
expected_rows,
file.physical(),
)
.map_err(|_| {
CoreExecutionError::ParquetContract(
"the retained file violates its signed physical envelope",
)
})?;
let mut reader = ExactParquetReader::from(ExactFileExtension {
file,
operation,
cancellation,
});
let metadata = reader.get_metadata(None).await?;
let observed_schema = parquet_to_arrow_schema(
metadata.file_metadata().schema_descr(),
metadata.file_metadata().key_value_metadata(),
)?;
if observed_schema.fields() != expected_schema.fields() {
return Err(CoreExecutionError::ParquetContract(
"the exact file schema differs from conversation-core",
));
}
let rows = metadata
.row_groups()
.iter()
.try_fold(0_u64, |held, group| {
let rows = u64::try_from(group.num_rows()).map_err(|_| {
CoreExecutionError::ParquetContract("a Parquet row count is negative")
})?;
if rows > 0 {
verify_source_statistics(group, &expected_schema, &expected_source)?;
}
held.checked_add(rows)
.ok_or(CoreExecutionError::ParquetContract(
"the Parquet row count overflowed",
))
})?;
if rows != expected_rows {
return Err(CoreExecutionError::ParquetContract(
"the exact file row count differs from its admitted segment",
));
}
Ok(())
}
fn verify_source_statistics(
group: &parquet::file::metadata::RowGroupMetaData,
schema: &SchemaRef,
source: &JournalSource,
) -> Result<(), CoreExecutionError> {
let partition = schema.index_of("partition")?;
let incarnation = schema.index_of("source_incarnation")?;
verify_constant_statistics(
group.column(partition).statistics(),
source.partition().as_str().as_bytes(),
)?;
verify_constant_statistics(
group.column(incarnation).statistics(),
source.incarnation().as_bytes(),
)
}
fn verify_constant_statistics(
statistics: Option<&parquet::file::statistics::Statistics>,
expected: &[u8],
) -> Result<(), CoreExecutionError> {
let statistics = statistics.ok_or(CoreExecutionError::ParquetContract(
"source lineage statistics are absent",
))?;
if !statistics.min_is_exact()
|| !statistics.max_is_exact()
|| statistics.null_count_opt() != Some(0)
|| statistics.min_bytes_opt() != Some(expected)
|| statistics.max_bytes_opt() != Some(expected)
{
return Err(CoreExecutionError::ParquetContract(
"source lineage statistics do not prove one exact source",
));
}
Ok(())
}
#[derive(Debug)]
pub(super) struct ExactParquetTable {
schema: SchemaRef,
files: Vec<PartitionedFile>,
factory: Arc<ExactParquetReaderFactory>,
}
impl ExactParquetTable {
pub(super) fn new(
schema: SchemaRef,
files: Vec<VerifiedCoreFile>,
operation: &Arc<CoreOperationContext>,
cancellation: &CancellationToken,
) -> Self {
let files = files
.into_iter()
.enumerate()
.map(|(index, file)| {
PartitionedFile::new(format!("exact-core/{index:08}.parquet"), file.byte_len())
.with_extension(ExactFileExtension {
file,
operation: Arc::clone(operation),
cancellation: cancellation.clone(),
})
})
.collect();
Self {
schema,
files,
factory: Arc::new(ExactParquetReaderFactory),
}
}
}
#[async_trait]
impl TableProvider for ExactParquetTable {
fn schema(&self) -> SchemaRef {
Arc::clone(&self.schema)
}
fn table_type(&self) -> TableType {
TableType::Base
}
async fn scan(
&self,
_state: &dyn Session,
projection: Option<&Vec<usize>>,
filters: &[Expr],
limit: Option<usize>,
) -> DataFusionResult<Arc<dyn ExecutionPlan>> {
if !filters.is_empty() {
return Err(DataFusionError::Plan(
"exact projection tables do not evaluate pushed filters".to_owned(),
));
}
if projection.is_some_and(|indices| {
indices
.iter()
.any(|index| *index >= self.schema.fields().len())
}) {
return Err(DataFusionError::Plan(
"exact projection table received an out-of-bound column".to_owned(),
));
}
let source = Arc::new(
ParquetSource::new(Arc::clone(&self.schema))
.with_parquet_file_reader_factory(self.factory.clone()),
);
let config = FileScanConfigBuilder::new(ObjectStoreUrl::local_filesystem(), source)
.with_file_groups(vec![FileGroup::new(self.files.clone())])
.with_preserve_order(true)
.with_limit(limit)
.with_projection_indices(projection.cloned())?
.build();
Ok(DataSourceExec::from_data_source(config))
}
fn supports_filters_pushdown(
&self,
filters: &[&Expr],
) -> DataFusionResult<Vec<TableProviderFilterPushDown>> {
Ok(vec![
TableProviderFilterPushDown::Unsupported;
filters.len()
])
}
}