use std::collections::HashMap;
use faucet_core::FaucetError;
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
use crate::credentials::DeltaCredentials;
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
pub struct DeltaConnection {
pub table_uri: String,
#[serde(default)]
pub credentials: DeltaCredentials,
#[serde(default)]
pub storage_options: HashMap<String, String>,
}
impl DeltaConnection {
pub fn merged_storage_options(&self) -> HashMap<String, String> {
let mut opts = self.storage_options.clone();
self.credentials.apply(&mut opts);
opts
}
pub fn redacted_uri(&self) -> String {
faucet_core::util::redact_uri_credentials(&self.table_uri)
}
pub fn register_handlers(&self) {
register_all_handlers();
}
fn table_url(&self) -> Result<url::Url, FaucetError> {
deltalake::ensure_table_uri(&self.table_uri).map_err(|e| {
FaucetError::Config(format!(
"delta: invalid table_uri '{}': {e}",
self.redacted_uri()
))
})
}
pub fn location_string(&self) -> Result<String, FaucetError> {
Ok(self.table_url()?.to_string())
}
pub async fn open(&self) -> Result<deltalake::DeltaTable, FaucetError> {
self.register_handlers();
let url = self.table_url()?;
deltalake::open_table_with_storage_options(url, self.merged_storage_options())
.await
.map_err(|e| self.map_open_error(e))
}
pub async fn open_optional(&self) -> Result<Option<deltalake::DeltaTable>, FaucetError> {
self.register_handlers();
let url = self.table_url()?;
match deltalake::open_table_with_storage_options(url, self.merged_storage_options()).await {
Ok(t) => Ok(Some(t)),
Err(e) if is_missing_table(&e) => Ok(None),
Err(e) => Err(self.map_open_error(e)),
}
}
pub async fn open_at_version(
&self,
version: u64,
) -> Result<deltalake::DeltaTable, FaucetError> {
let mut table = self.open().await?;
table
.load_version(version)
.await
.map_err(|e| self.map_open_error(e))?;
Ok(table)
}
pub async fn open_at_timestamp(
&self,
timestamp: &str,
) -> Result<deltalake::DeltaTable, FaucetError> {
let ts = chrono::DateTime::parse_from_rfc3339(timestamp)
.map_err(|e| {
FaucetError::Config(format!(
"delta: invalid `timestamp` '{timestamp}' (expected RFC 3339): {e}"
))
})?
.with_timezone(&chrono::Utc);
let mut table = self.open().await?;
table
.load_with_datetime(ts)
.await
.map_err(|e| self.map_open_error(e))?;
Ok(table)
}
fn map_open_error(&self, e: deltalake::DeltaTableError) -> FaucetError {
FaucetError::Source(format!(
"delta: failed to open table '{}': {e}",
self.redacted_uri()
))
}
}
fn is_missing_table(e: &deltalake::DeltaTableError) -> bool {
if matches!(e, deltalake::DeltaTableError::NotATable(_)) {
return true;
}
let msg = e.to_string().to_lowercase();
msg.contains("not found") || msg.contains("no such file") || msg.contains("does not exist")
}
fn register_all_handlers() {
#[cfg(feature = "s3")]
{
use std::sync::Once;
static S3: Once = Once::new();
S3.call_once(|| deltalake::aws::register_handlers(None));
}
#[cfg(feature = "azure")]
{
use std::sync::Once;
static AZURE: Once = Once::new();
AZURE.call_once(|| deltalake::azure::register_handlers(None));
}
#[cfg(feature = "gcs")]
{
use std::sync::Once;
static GCS: Once = Once::new();
GCS.call_once(|| deltalake::gcp::register_handlers(None));
}
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
fn conn(uri: &str) -> DeltaConnection {
DeltaConnection {
table_uri: uri.into(),
credentials: DeltaCredentials::default(),
storage_options: HashMap::new(),
}
}
#[test]
fn flatten_shape_round_trips() {
let v = json!({
"table_uri": "s3://bucket/tbl",
"credentials": { "type": "aws", "config": { "region": "us-east-1" } },
"storage_options": { "AWS_ALLOW_HTTP": "true" }
});
let c: DeltaConnection = serde_json::from_value(v).unwrap();
assert_eq!(c.table_uri, "s3://bucket/tbl");
let opts = c.merged_storage_options();
assert_eq!(opts["AWS_REGION"], "us-east-1");
assert_eq!(opts["AWS_ALLOW_HTTP"], "true");
}
#[test]
fn explicit_storage_option_overrides_credential() {
let v = json!({
"table_uri": "s3://bucket/tbl",
"credentials": { "type": "aws", "config": { "region": "us-east-1" } },
"storage_options": { "AWS_REGION": "eu-west-1" }
});
let c: DeltaConnection = serde_json::from_value(v).unwrap();
assert_eq!(c.merged_storage_options()["AWS_REGION"], "eu-west-1");
}
#[test]
fn defaults_omit_credentials_and_options() {
let c: DeltaConnection =
serde_json::from_value(json!({ "table_uri": "file:///tmp/t" })).unwrap();
assert_eq!(c.credentials, DeltaCredentials::Default);
assert!(c.merged_storage_options().is_empty());
}
#[test]
fn redacted_uri_passes_plain_uri_through() {
assert_eq!(conn("s3://bucket/tbl").redacted_uri(), "s3://bucket/tbl");
}
#[test]
fn register_handlers_is_idempotent() {
let c = conn("file:///tmp/t");
c.register_handlers();
c.register_handlers();
}
}