use std::sync::Arc;
use arrow::array::Int64Array;
use arrow::datatypes::{DataType, Field, Schema};
use arrow::record_batch::RecordBatch;
use deltalake::kernel::StructType;
use deltalake::kernel::engine::arrow_conversion::TryIntoKernel;
use deltalake::operations::create::CreateBuilder;
use deltalake::writer::{DeltaWriter, RecordBatchWriter};
use faucet_conformance::{
assert_batch_size_zero_single_page, assert_bounded_memory, assert_config_schema_valid_value,
assert_connector_name_nonempty, assert_errors_not_panics, assert_preflight_check_wellformed,
};
use faucet_core::Source as _;
use faucet_source_delta::{DeltaSource, DeltaSourceConfig};
fn table_uri(dir: &tempfile::TempDir, name: &str) -> String {
dir.path().join(name).to_string_lossy().into_owned()
}
async fn seed(uri: &str, n: i64) {
let arrow_schema = Arc::new(Schema::new(vec![Field::new("id", DataType::Int64, false)]));
let ids = Int64Array::from((0..n).collect::<Vec<i64>>());
let batch =
RecordBatch::try_new(arrow_schema.clone(), vec![Arc::new(ids)]).expect("seed batch");
let delta_schema: StructType = arrow_schema
.as_ref()
.try_into_kernel()
.expect("arrow schema → delta kernel");
let mut table = CreateBuilder::new()
.with_location(uri)
.with_columns(delta_schema.fields().cloned())
.await
.expect("create table");
let mut writer = RecordBatchWriter::for_table(&table).expect("record-batch writer");
writer.write(batch).await.expect("write batch");
writer
.flush_and_commit(&mut table)
.await
.expect("commit batch");
}
#[test]
fn conformance_config_schema_valid() {
let schema = serde_json::to_value(schemars::schema_for!(DeltaSourceConfig)).unwrap();
assert_config_schema_valid_value(&schema, "delta");
}
#[tokio::test(flavor = "multi_thread")]
async fn conformance_bounded_memory() {
let dir = tempfile::tempdir().unwrap();
let uri = table_uri(&dir, "events");
seed(&uri, 250).await;
let mut cfg = DeltaSourceConfig::new(&uri);
cfg.batch_size = 50;
let source = DeltaSource::new(cfg).await.expect("source");
assert_bounded_memory(&source, 50, 250).await;
let all = source.fetch_all().await.expect("fetch_all");
assert_eq!(all.len(), 250);
}
#[tokio::test(flavor = "multi_thread")]
async fn conformance_errors_not_panics() {
let dir = tempfile::tempdir().unwrap();
let uri = table_uri(&dir, "not-a-table");
if let Ok(source) = DeltaSource::new(DeltaSourceConfig::new(&uri)).await {
assert_errors_not_panics(&source).await;
}
}
#[tokio::test(flavor = "multi_thread")]
async fn conformance_connector_name_nonempty() {
let dir = tempfile::tempdir().unwrap();
let uri = table_uri(&dir, "name");
seed(&uri, 1).await;
let source = DeltaSource::new(DeltaSourceConfig::new(&uri))
.await
.expect("source");
assert_connector_name_nonempty(&source);
}
#[tokio::test(flavor = "multi_thread")]
async fn conformance_batch_size_zero_single_page() {
let dir = tempfile::tempdir().unwrap();
let uri = table_uri(&dir, "single_page");
seed(&uri, 6).await;
let mut cfg = DeltaSourceConfig::new(&uri);
cfg.batch_size = 0;
let source = DeltaSource::new(cfg).await.expect("source");
assert_batch_size_zero_single_page(&source).await;
}
#[tokio::test(flavor = "multi_thread")]
async fn conformance_preflight_check_wellformed() {
let dir = tempfile::tempdir().unwrap();
let uri = table_uri(&dir, "preflight");
seed(&uri, 6).await;
let source = DeltaSource::new(DeltaSourceConfig::new(&uri))
.await
.expect("source");
assert_preflight_check_wellformed(&source, &faucet_core::check::CheckContext::default()).await;
}