use bytes::Bytes;
use futures::StreamExt;
use minio::s3::MinioClient;
use minio::s3::creds::StaticProvider;
use minio::s3::error::{Error as S3Error, S3ServerError};
use minio::s3::http::BaseUrl;
use minio::s3::minio_error_response::MinioErrorCode;
use minio::s3::response_traits::HasEtagFromHeaders;
use minio::s3::segmented_bytes::SegmentedBytes;
use minio::s3::types::{S3Api, ToStream};
use crate::dataops::store::BucketConfig;
pub struct ObjectEntry {
pub key: String,
pub size: u64,
pub is_prefix: bool,
}
pub(crate) enum UploadOutcome {
Uploaded,
Unchanged,
}
fn build_client(bucket_config: &BucketConfig, secret_key: &str) -> Result<MinioClient, String> {
let base_url: BaseUrl = bucket_config
.endpoint
.parse()
.map_err(|err| format!("invalid endpoint '{}': {err}", bucket_config.endpoint))?;
let provider = StaticProvider::new(&bucket_config.access_key_id, secret_key, None);
MinioClient::new(base_url, Some(provider), None, None).map_err(|err| {
format!(
"failed to create S3 client for '{}': {err}",
bucket_config.alias
)
})
}
fn format_error(err: &S3Error) -> String {
if let S3Error::S3Server(S3ServerError::S3Error(response)) = err {
let code = response.code();
return match response.message() {
Some(message) => format!("{code:?}: {message}"),
None => format!("{code:?}"),
};
}
err.to_string()
}
pub async fn list_buckets(
bucket_config: &BucketConfig,
secret_key: &str,
) -> Result<Vec<String>, String> {
let client = build_client(bucket_config, secret_key)?;
let resp = client
.list_buckets()
.build()
.send()
.await
.map_err(|err| format!("failed to list buckets: {}", format_error(&err)))?;
let buckets = resp
.buckets()
.map_err(|err| format!("failed to parse bucket list: {err}"))?;
Ok(buckets
.into_iter()
.map(|bucket| bucket.name.to_string())
.collect())
}
pub(crate) async fn bucket_exists(
bucket_config: &BucketConfig,
secret_key: &str,
) -> Result<bool, String> {
let client = build_client(bucket_config, secret_key)?;
let resp = client
.bucket_exists(bucket_config.bucket.as_str())
.map_err(|err| format!("invalid bucket name '{}': {err}", bucket_config.bucket))?
.build()
.send()
.await
.map_err(|err| format!("failed to check bucket: {}", format_error(&err)))?;
Ok(resp.exists())
}
pub async fn list_objects(
bucket_config: &BucketConfig,
secret_key: &str,
prefix: &str,
recursive: bool,
) -> Result<Vec<ObjectEntry>, String> {
let client = build_client(bucket_config, secret_key)?;
let prefix = Some(prefix.to_string());
let list = if recursive {
client
.list_objects(bucket_config.bucket.as_str())
.map_err(|err| format!("invalid bucket name '{}': {err}", bucket_config.bucket))?
.prefix(prefix)
.recursive(true)
.build()
} else {
client
.list_objects(bucket_config.bucket.as_str())
.map_err(|err| format!("invalid bucket name '{}': {err}", bucket_config.bucket))?
.prefix(prefix)
.delimiter(Some("/".to_string()))
.build()
};
let mut stream = list.to_stream().await;
let mut entries = Vec::new();
while let Some(page) = stream.next().await {
let page = page.map_err(|err| format!("failed to list objects: {}", format_error(&err)))?;
for item in page.contents {
entries.push(ObjectEntry {
key: item.name,
size: item.size.unwrap_or(0),
is_prefix: item.is_prefix,
});
}
}
Ok(entries)
}
pub async fn get_object(
bucket_config: &BucketConfig,
secret_key: &str,
key: &str,
) -> Result<Vec<u8>, String> {
let client = build_client(bucket_config, secret_key)?;
let resp = client
.get_object(bucket_config.bucket.as_str(), key)
.map_err(|err| format!("invalid object key '{key}': {err}"))?
.build()
.send()
.await
.map_err(|err| format!("failed to download '{key}': {}", format_error(&err)))?;
let content = resp
.content()
.map_err(|err| format!("failed to read '{key}': {err}"))?
.to_segmented_bytes()
.await
.map_err(|err| format!("failed to read '{key}': {err}"))?;
Ok(content.to_bytes().to_vec())
}
pub(crate) async fn upload_if_changed(
bucket_config: &BucketConfig,
secret_key: &str,
key: &str,
data: Vec<u8>,
) -> Result<UploadOutcome, String> {
let client = build_client(bucket_config, secret_key)?;
let existing_etag = match client
.stat_object(bucket_config.bucket.as_str(), key)
.map_err(|err| format!("invalid object key '{key}': {err}"))?
.build()
.send()
.await
{
Ok(resp) => Some(
resp.etag()
.map_err(|err| format!("failed to read ETag for '{key}': {err}"))?
.to_string(),
),
Err(S3Error::S3Server(S3ServerError::S3Error(ref response)))
if response.code() == MinioErrorCode::NoSuchKey =>
{
None
}
Err(err) => {
return Err(format!("failed to check '{key}': {}", format_error(&err)));
}
};
let local_hash = format!("{:x}", md5::compute(&data));
if let Some(existing) = existing_etag {
if existing == local_hash {
return Ok(UploadOutcome::Unchanged);
}
println!("note: '{key}' changed since last upload, new version created");
}
let bytes = SegmentedBytes::from(Bytes::from(data));
client
.put_object(bucket_config.bucket.as_str(), key, bytes)
.map_err(|err| format!("invalid object key '{key}': {err}"))?
.build()
.send()
.await
.map_err(|err| format!("failed to upload '{key}': {}", format_error(&err)))?;
Ok(UploadOutcome::Uploaded)
}