use std::sync::Arc;
use bytes::Bytes;
use deltalake_core::logstore::*;
use deltalake_core::{
DeltaResult, kernel::Version, kernel::transaction::TransactionError, logstore::ObjectStoreRef,
};
use object_store::{Error as ObjectStoreError, ObjectStore};
use url::Url;
use uuid::Uuid;
pub fn default_s3_logstore(
store: ObjectStoreRef,
root_store: ObjectStoreRef,
location: &Url,
options: &StorageConfig,
) -> Arc<dyn LogStore> {
Arc::new(S3LogStore::new(
store,
root_store,
LogStoreConfig::new(location, options.clone()),
))
}
#[derive(Debug, Clone)]
pub struct S3LogStore {
prefixed_store: ObjectStoreRef,
root_store: ObjectStoreRef,
config: LogStoreConfig,
}
impl S3LogStore {
pub fn new(
prefixed_store: ObjectStoreRef,
root_store: ObjectStoreRef,
config: LogStoreConfig,
) -> Self {
Self {
prefixed_store,
root_store,
config,
}
}
}
#[async_trait::async_trait]
impl LogStore for S3LogStore {
fn name(&self) -> String {
"S3LogStore".into()
}
async fn read_commit_entry(&self, version: Version) -> DeltaResult<Option<Bytes>> {
read_commit_entry(self.object_store(None).as_ref(), version).await
}
async fn write_commit_entry(
&self,
version: Version,
commit_or_bytes: CommitOrBytes,
_operation_id: Uuid,
) -> Result<(), TransactionError> {
match commit_or_bytes {
CommitOrBytes::TmpCommit(tmp_commit) => {
Ok(
write_commit_entry(self.object_store(None).as_ref(), version, &tmp_commit)
.await?,
)
}
_ => unreachable!(), }
.map_err(|err| -> TransactionError {
match err {
ObjectStoreError::AlreadyExists { .. } => {
TransactionError::VersionAlreadyExists(version)
}
_ => TransactionError::from(err),
}
})?;
Ok(())
}
async fn abort_commit_entry(
&self,
version: Version,
commit_or_bytes: CommitOrBytes,
_operation_id: Uuid,
) -> Result<(), TransactionError> {
match &commit_or_bytes {
CommitOrBytes::TmpCommit(tmp_commit) => {
abort_commit_entry(self.object_store(None).as_ref(), version, tmp_commit).await
}
_ => unreachable!(), }
}
async fn get_latest_version(&self, current_version: Version) -> DeltaResult<Version> {
get_latest_version(self, current_version).await
}
fn object_store(&self, _operation_id: Option<Uuid>) -> Arc<dyn ObjectStore> {
self.prefixed_store.clone()
}
fn root_object_store(&self, _operation_id: Option<Uuid>) -> Arc<dyn ObjectStore> {
self.root_store.clone()
}
fn config(&self) -> &LogStoreConfig {
&self.config
}
}