use super::Connector;
use crate::profiles::Profile;
use anyhow::{Context, Result, anyhow};
use aws_config::BehaviorVersion;
use aws_sdk_s3::Client as S3Client;
use aws_sdk_s3::config::Credentials;
use std::io::{Cursor, Read};
pub struct S3Connector {
client: S3Client,
bucket: String,
}
impl S3Connector {
pub async fn from_profile_and_url(profile: &Profile, url: &url::Url) -> Result<Self> {
let bucket = url
.host_str()
.ok_or_else(|| anyhow!("Invalid S3 URL: missing bucket name"))?
.to_string();
let region = profile
.region
.clone()
.unwrap_or_else(|| "us-east-1".to_string());
let mut config_loader =
aws_config::defaults(BehaviorVersion::latest()).region(aws_config::Region::new(region));
if let Some(endpoint) = &profile.endpoint {
config_loader = config_loader.endpoint_url(endpoint);
}
let base_config = config_loader.load().await;
let mut s3_config = aws_sdk_s3::config::Builder::from(&base_config);
if let (Some(access_key), Some(secret_key)) = (&profile.access_key, &profile.secret_key) {
if !access_key.is_empty() && !secret_key.is_empty() {
let creds = Credentials::new(
access_key.clone(),
secret_key.clone(),
None,
None,
"profile",
);
s3_config = s3_config.credentials_provider(creds);
}
}
if profile.path_style.unwrap_or(false) {
s3_config = s3_config.force_path_style(true);
}
let client = S3Client::from_conf(s3_config.build());
Ok(S3Connector { client, bucket })
}
pub async fn put_object_from_url(&self, s3_url: &str, data: &[u8]) -> Result<()> {
let url = url::Url::parse(s3_url)?;
let bucket = url
.host_str()
.ok_or_else(|| anyhow!("Invalid S3 URL: missing bucket"))?;
let key = url.path().trim_start_matches('/');
use aws_sdk_s3::primitives::ByteStream;
self.client
.put_object()
.bucket(bucket)
.key(key)
.body(ByteStream::from(data.to_vec()))
.send()
.await
.context("Failed to upload S3 object")?;
Ok(())
}
fn parse_s3_path(&self, path: &str) -> Result<String> {
if path.starts_with("s3://") {
let url = url::Url::parse(path)?;
Ok(url.path().trim_start_matches('/').to_string())
} else {
Ok(path.to_string())
}
}
}
#[async_trait::async_trait]
impl Connector for S3Connector {
async fn fetch(&self, location: &str) -> Result<Box<dyn Read>> {
let key = self.parse_s3_path(location)?;
let resp = self
.client
.get_object()
.bucket(&self.bucket)
.key(&key)
.send()
.await
.context("Failed to fetch S3 object")?;
let data = resp
.body
.collect()
.await
.context("Failed to read S3 object body")?
.into_bytes();
Ok(Box::new(Cursor::new(data)))
}
}