mirror of
https://tangled.org/tranquil.farm/tranquil-pds
synced 2026-09-28 05:04:14 +00:00
refactor: drop aws crates for s3
We were just using them for S3 support. This PR switches us over to the s3 crate which works the same way, with some minor differences. The main implication of this change is that we're no longer doing all the AWS stuff, like checking aws profile in home directory, or automatic discovery of metadata. This change is covered by our brand new s3 test suite, but I also did the old school `just test-full` with minio.
This commit is contained in:
Generated
+392
-701
File diff suppressed because it is too large
Load Diff
+1
-3
@@ -63,9 +63,6 @@ backon = "1"
|
||||
bincode = { version = "2", features = ["serde"] }
|
||||
anyhow = "1.0"
|
||||
async-trait = "0.1"
|
||||
aws-config = { version = "1.8", default-features = false, features = ["rt-tokio", "credentials-process", "sso"] }
|
||||
aws-sdk-s3 = { version = "1.118", default-features = false, features = ["sigv4a", "rt-tokio"] }
|
||||
aws-smithy-http-client = { version = "1.1", features = ["rustls-ring"] }
|
||||
axum = { version = "0.8", features = ["ws", "macros"] }
|
||||
base32 = "0.5"
|
||||
base64 = "0.22"
|
||||
@@ -116,6 +113,7 @@ regex = "1"
|
||||
rsa = "0.9"
|
||||
rustls = { version = "0.23", default-features = false, features = ["ring", "std", "tls12", "logging"] }
|
||||
rustls-pemfile = "2"
|
||||
s3 = { version = "0.1.36", features = ["multipart"] }
|
||||
s3s = "0.17"
|
||||
s3s-fs = "0.17"
|
||||
secrecy = { version = "0.10", features = ["serde"] }
|
||||
|
||||
@@ -755,6 +755,7 @@ impl fmt::Display for RepoBackend {
|
||||
#[config(layer_attr(serde(deny_unknown_fields)))]
|
||||
pub struct StorageConfig {
|
||||
/// Storage backend: `filesystem` or `s3`.
|
||||
/// If you set `s3` you must also provide `AWS_ACCESS_KEY_ID` and `AWS_SECRET_ACCESS_KEY`.
|
||||
#[config(env = "BLOB_STORAGE_BACKEND", default = "filesystem")]
|
||||
pub backend: String,
|
||||
|
||||
|
||||
@@ -100,7 +100,5 @@ tracing-subscriber = { workspace = true, features = ["env-filter"] }
|
||||
wiremock = { workspace = true }
|
||||
hex = { workspace = true }
|
||||
tokio-tungstenite = { workspace = true }
|
||||
aws-config = { workspace = true }
|
||||
aws-sdk-s3 = { workspace = true }
|
||||
aws-smithy-http-client = { workspace = true }
|
||||
s3 = { workspace = true }
|
||||
redis = { workspace = true }
|
||||
|
||||
@@ -1,9 +1,3 @@
|
||||
#[cfg(all(not(feature = "external-infra"), feature = "s3"))]
|
||||
use aws_config::BehaviorVersion;
|
||||
#[cfg(all(not(feature = "external-infra"), feature = "s3"))]
|
||||
use aws_sdk_s3::Client as S3Client;
|
||||
#[cfg(all(not(feature = "external-infra"), feature = "s3"))]
|
||||
use aws_sdk_s3::config::Credentials;
|
||||
use chrono::Utc;
|
||||
use reqwest::{Client, StatusCode, header};
|
||||
use serde_json::{Value, json};
|
||||
@@ -59,7 +53,7 @@ pub struct ServerInstance {
|
||||
#[cfg(all(not(feature = "external-infra"), feature = "s3"))]
|
||||
use testcontainers::GenericImage;
|
||||
#[cfg(all(not(feature = "external-infra"), feature = "s3"))]
|
||||
use testcontainers::core::ContainerPort;
|
||||
use testcontainers::core::{ContainerPort, WaitFor};
|
||||
#[cfg(not(feature = "external-infra"))]
|
||||
use testcontainers::{ContainerAsync, ImageExt, runners::AsyncRunner};
|
||||
#[cfg(not(feature = "external-infra"))]
|
||||
@@ -243,6 +237,7 @@ async fn setup_with_testcontainers() -> String {
|
||||
#[cfg(all(not(feature = "external-infra"), feature = "s3"))]
|
||||
async fn setup_with_testcontainers() -> String {
|
||||
let s3_container = GenericImage::new("cgr.dev/chainguard/minio", "latest")
|
||||
.with_wait_for(WaitFor::message_on_stderr("API: http"))
|
||||
.with_exposed_port(ContainerPort::Tcp(9000))
|
||||
.with_env_var("MINIO_ROOT_USER", "minioadmin")
|
||||
.with_env_var("MINIO_ROOT_PASSWORD", "minioadmin")
|
||||
@@ -268,35 +263,25 @@ async fn setup_with_testcontainers() -> String {
|
||||
std::env::set_var("SKIP_IMPORT_VERIFICATION", "true");
|
||||
std::env::set_var("PLC_DIRECTORY_URL", &plc_url);
|
||||
}
|
||||
let sdk_config = aws_config::defaults(BehaviorVersion::latest())
|
||||
let s3_client = s3::Client::builder(&s3_endpoint)
|
||||
.unwrap()
|
||||
.region("us-east-1")
|
||||
.http_client(
|
||||
aws_smithy_http_client::Builder::new()
|
||||
.tls_provider(aws_smithy_http_client::tls::Provider::Rustls(
|
||||
aws_smithy_http_client::tls::rustls_provider::CryptoMode::Ring,
|
||||
))
|
||||
.build_https(),
|
||||
)
|
||||
.endpoint_url(&s3_endpoint)
|
||||
.credentials_provider(Credentials::new(
|
||||
"minioadmin",
|
||||
"minioadmin",
|
||||
None,
|
||||
None,
|
||||
"test",
|
||||
))
|
||||
.load()
|
||||
.await;
|
||||
let s3_config = aws_sdk_s3::config::Builder::from(&sdk_config)
|
||||
.force_path_style(true)
|
||||
.build();
|
||||
let s3_client = S3Client::from_conf(s3_config);
|
||||
let _ = s3_client.create_bucket().bucket("test-bucket").send().await;
|
||||
let _ = s3_client
|
||||
.create_bucket()
|
||||
.bucket("test-backups")
|
||||
.addressing_style(s3::AddressingStyle::Path)
|
||||
.auth(s3::Auth::from_env().unwrap())
|
||||
.build()
|
||||
.unwrap();
|
||||
s3_client
|
||||
.buckets()
|
||||
.create("test-bucket")
|
||||
.send()
|
||||
.await;
|
||||
.await
|
||||
.unwrap();
|
||||
s3_client
|
||||
.buckets()
|
||||
.create("test-backups")
|
||||
.send()
|
||||
.await
|
||||
.unwrap();
|
||||
register_mock_appview().await;
|
||||
S3_CONTAINER.set(s3_container).ok();
|
||||
let container = Postgres::default()
|
||||
|
||||
@@ -5,18 +5,16 @@ edition.workspace = true
|
||||
license.workspace = true
|
||||
|
||||
[features]
|
||||
s3 = ["dep:aws-config", "dep:aws-sdk-s3", "dep:aws-smithy-http-client"]
|
||||
s3 = ["dep:s3"]
|
||||
|
||||
[dependencies]
|
||||
tranquil-config = { workspace = true }
|
||||
tranquil-infra = { workspace = true }
|
||||
|
||||
async-trait = { workspace = true }
|
||||
aws-config = { workspace = true, optional = true }
|
||||
aws-sdk-s3 = { workspace = true, optional = true }
|
||||
aws-smithy-http-client = { workspace = true, optional = true }
|
||||
bytes = { workspace = true }
|
||||
futures = { workspace = true }
|
||||
s3 = { workspace = true, optional = true }
|
||||
sha2 = { workspace = true }
|
||||
tokio = { workspace = true }
|
||||
tracing = { workspace = true }
|
||||
|
||||
@@ -102,12 +102,8 @@ fn map_io_not_found(key: &str) -> impl FnOnce(std::io::Error) -> StorageError +
|
||||
#[cfg(feature = "s3")]
|
||||
mod s3 {
|
||||
use super::*;
|
||||
use aws_config::BehaviorVersion;
|
||||
use aws_sdk_s3::Client;
|
||||
use aws_sdk_s3::config::Region;
|
||||
use aws_sdk_s3::primitives::ByteStream;
|
||||
use aws_sdk_s3::types::CompletedMultipartUpload;
|
||||
use aws_sdk_s3::types::CompletedPart;
|
||||
use ::s3::types::CompletedPart;
|
||||
use ::s3::{AddressingStyle, Auth, Client};
|
||||
|
||||
const MIN_PART_SIZE: usize = 5 * 1024 * 1024;
|
||||
|
||||
@@ -120,7 +116,7 @@ mod s3 {
|
||||
impl S3BlobStorage {
|
||||
pub async fn new(bucket: &str, endpoint: Option<&str>, path: &str) -> Self {
|
||||
Self {
|
||||
client: create_s3_client(endpoint).await,
|
||||
client: create_s3_client(endpoint),
|
||||
bucket: bucket.to_string(),
|
||||
path: path
|
||||
.trim_start_matches("/")
|
||||
@@ -138,28 +134,26 @@ mod s3 {
|
||||
}
|
||||
}
|
||||
|
||||
async fn create_s3_client(endpoint: Option<&str>) -> Client {
|
||||
let http_client = aws_smithy_http_client::Builder::new()
|
||||
.tls_provider(aws_smithy_http_client::tls::Provider::Rustls(
|
||||
aws_smithy_http_client::tls::rustls_provider::CryptoMode::Ring,
|
||||
))
|
||||
.build_https();
|
||||
|
||||
let config = aws_config::defaults(BehaviorVersion::latest())
|
||||
.http_client(http_client)
|
||||
.load()
|
||||
.await;
|
||||
|
||||
let region = config
|
||||
.region()
|
||||
.cloned()
|
||||
.unwrap_or_else(|| Region::from_static("us-east-1"));
|
||||
let builder = aws_sdk_s3::config::Builder::from(&config).region(region);
|
||||
let builder = match endpoint {
|
||||
Some(endpoint) => builder.endpoint_url(endpoint).force_path_style(true),
|
||||
None => builder,
|
||||
fn create_s3_client(endpoint: Option<&str>) -> Client {
|
||||
let region = std::env::var("AWS_REGION").unwrap_or_else(|_| "us-east-1".to_string());
|
||||
let (url, style) = match endpoint {
|
||||
Some(endpoint) => (endpoint.to_string(), AddressingStyle::Path),
|
||||
None => (
|
||||
format!("https://s3.{region}.amazonaws.com"),
|
||||
AddressingStyle::Auto,
|
||||
),
|
||||
};
|
||||
Client::from_conf(builder.build())
|
||||
let auth =
|
||||
Auth::from_env().expect("AWS_ACCESS_KEY_ID and AWS_SECRET_ACCESS_KEY must be set");
|
||||
Client::builder(url)
|
||||
.expect("storage.s3_endpoint (S3_ENDPOINT) must be a valid URL")
|
||||
.region(region)
|
||||
.addressing_style(style)
|
||||
.auth(auth)
|
||||
.base_retry_delay(std::time::Duration::from_secs(1))
|
||||
.max_retry_delay(std::time::Duration::from_secs(20))
|
||||
.build()
|
||||
.expect("failed to build S3 client")
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
@@ -171,10 +165,9 @@ mod s3 {
|
||||
async fn put_bytes(&self, key: &str, data: Bytes) -> Result<(), StorageError> {
|
||||
let path = self.resolve_path(key);
|
||||
self.client
|
||||
.put_object()
|
||||
.bucket(&self.bucket)
|
||||
.key(&path)
|
||||
.body(ByteStream::from(data))
|
||||
.objects()
|
||||
.put(&self.bucket, path)
|
||||
.body_bytes(data)
|
||||
.send()
|
||||
.await
|
||||
.map_err(|e| StorageError::Backend(e.to_string()))?;
|
||||
@@ -190,46 +183,39 @@ mod s3 {
|
||||
let path = self.resolve_path(key);
|
||||
let resp = self
|
||||
.client
|
||||
.get_object()
|
||||
.bucket(&self.bucket)
|
||||
.key(&path)
|
||||
.objects()
|
||||
.get(&self.bucket, &path)
|
||||
.send()
|
||||
.await
|
||||
.map_err(|e| StorageError::Backend(e.to_string()))?;
|
||||
|
||||
resp.body
|
||||
.collect()
|
||||
resp.bytes()
|
||||
.await
|
||||
.map(|agg| agg.into_bytes())
|
||||
.map_err(|e| StorageError::Backend(e.to_string()))
|
||||
}
|
||||
|
||||
async fn get_head(&self, key: &str, size: usize) -> Result<Bytes, StorageError> {
|
||||
let path = self.resolve_path(key);
|
||||
let range = format!("bytes=0-{}", size.saturating_sub(1));
|
||||
let resp = self
|
||||
.client
|
||||
.get_object()
|
||||
.bucket(&self.bucket)
|
||||
.key(&path)
|
||||
.range(range)
|
||||
.objects()
|
||||
.get(&self.bucket, &path)
|
||||
.range_bytes(0, size.saturating_sub(1) as u64)
|
||||
.map_err(|e| StorageError::Backend(e.to_string()))?
|
||||
.send()
|
||||
.await
|
||||
.map_err(|e| StorageError::Backend(e.to_string()))?;
|
||||
|
||||
resp.body
|
||||
.collect()
|
||||
resp.bytes()
|
||||
.await
|
||||
.map(|agg| agg.into_bytes())
|
||||
.map_err(|e| StorageError::Backend(e.to_string()))
|
||||
}
|
||||
|
||||
async fn delete(&self, key: &str) -> Result<(), StorageError> {
|
||||
let path = self.resolve_path(key);
|
||||
self.client
|
||||
.delete_object()
|
||||
.bucket(&self.bucket)
|
||||
.key(&path)
|
||||
.objects()
|
||||
.delete(&self.bucket, &path)
|
||||
.send()
|
||||
.await
|
||||
.map_err(|e| StorageError::Backend(e.to_string()))?;
|
||||
@@ -245,27 +231,22 @@ mod s3 {
|
||||
use futures::StreamExt;
|
||||
|
||||
let path = self.resolve_path(key);
|
||||
let create_resp = self
|
||||
let upload_id = self
|
||||
.client
|
||||
.create_multipart_upload()
|
||||
.bucket(&self.bucket)
|
||||
.key(&path)
|
||||
.objects()
|
||||
.create_multipart_upload(&self.bucket, &path)
|
||||
.send()
|
||||
.await
|
||||
.map_err(|e| {
|
||||
StorageError::Backend(format!("Failed to create multipart upload: {}", e))
|
||||
})?;
|
||||
|
||||
let upload_id = create_resp
|
||||
.upload_id()
|
||||
.ok_or_else(|| StorageError::Backend("No upload ID returned".to_string()))?
|
||||
.to_string();
|
||||
})?
|
||||
.upload_id;
|
||||
|
||||
let upload_part = |client: &Client,
|
||||
bucket: &str,
|
||||
key: &str,
|
||||
upload_id: &str,
|
||||
part_num: i32,
|
||||
part_num: u32,
|
||||
data: Vec<u8>|
|
||||
-> std::pin::Pin<
|
||||
Box<dyn std::future::Future<Output = Result<CompletedPart, StorageError>> + Send>,
|
||||
@@ -276,36 +257,28 @@ mod s3 {
|
||||
let upload_id = upload_id.to_string();
|
||||
Box::pin(async move {
|
||||
let resp = client
|
||||
.upload_part()
|
||||
.bucket(&bucket)
|
||||
.key(&path)
|
||||
.upload_id(&upload_id)
|
||||
.part_number(part_num)
|
||||
.body(ByteStream::from(data))
|
||||
.objects()
|
||||
.upload_part(&bucket, &path, &upload_id, part_num)
|
||||
.body_bytes(data)
|
||||
.send()
|
||||
.await
|
||||
.map_err(|e| {
|
||||
StorageError::Backend(format!("Failed to upload part: {}", e))
|
||||
})?;
|
||||
|
||||
let etag = resp
|
||||
.e_tag()
|
||||
.ok_or_else(|| {
|
||||
StorageError::Backend("No ETag returned for part".to_string())
|
||||
})?
|
||||
.to_string();
|
||||
let etag = resp.etag.ok_or_else(|| {
|
||||
StorageError::Backend("No ETag returned for part".to_string())
|
||||
})?;
|
||||
|
||||
Ok(CompletedPart::builder()
|
||||
.part_number(part_num)
|
||||
.e_tag(etag)
|
||||
.build())
|
||||
CompletedPart::new(part_num, etag)
|
||||
.map_err(|e| StorageError::Backend(e.to_string()))
|
||||
})
|
||||
};
|
||||
|
||||
struct UploadState {
|
||||
hasher: Sha256,
|
||||
total_size: u64,
|
||||
part_number: i32,
|
||||
part_number: u32,
|
||||
completed_parts: Vec<CompletedPart>,
|
||||
buffer: Vec<u8>,
|
||||
}
|
||||
@@ -321,10 +294,8 @@ mod s3 {
|
||||
let abort_upload = || async {
|
||||
let _ = self
|
||||
.client
|
||||
.abort_multipart_upload()
|
||||
.bucket(&self.bucket)
|
||||
.key(path)
|
||||
.upload_id(&upload_id)
|
||||
.objects()
|
||||
.abort_multipart_upload(&self.bucket, path, &upload_id)
|
||||
.send()
|
||||
.await;
|
||||
};
|
||||
@@ -378,17 +349,12 @@ mod s3 {
|
||||
return Err(StorageError::Other("Empty upload".to_string()));
|
||||
}
|
||||
|
||||
let completed_upload = CompletedMultipartUpload::builder()
|
||||
.set_parts(Some(state.completed_parts))
|
||||
.build();
|
||||
|
||||
let path = self.resolve_path(key);
|
||||
self.client
|
||||
.complete_multipart_upload()
|
||||
.bucket(&self.bucket)
|
||||
.key(&path)
|
||||
.upload_id(&upload_id)
|
||||
.multipart_upload(completed_upload)
|
||||
.objects()
|
||||
.complete_multipart_upload(&self.bucket, &path, &upload_id)
|
||||
.parts(state.completed_parts)
|
||||
.map_err(|e| StorageError::Backend(e.to_string()))?
|
||||
.send()
|
||||
.await
|
||||
.map_err(|e| {
|
||||
@@ -411,14 +377,11 @@ mod s3 {
|
||||
|
||||
async fn copy(&self, src_key: &str, dst_key: &str) -> Result<(), StorageError> {
|
||||
let src_path = self.resolve_path(src_key);
|
||||
let copy_source = format!("{}/{}", self.bucket, &src_path);
|
||||
let dst_path = self.resolve_path(dst_key);
|
||||
|
||||
self.client
|
||||
.copy_object()
|
||||
.bucket(&self.bucket)
|
||||
.copy_source(copy_source)
|
||||
.key(&dst_path)
|
||||
.objects()
|
||||
.copy(&self.bucket, &src_path, &self.bucket, &dst_path)
|
||||
.send()
|
||||
.await
|
||||
.map_err(|e| StorageError::Backend(format!("Failed to copy object: {}", e)))?;
|
||||
|
||||
@@ -241,6 +241,7 @@
|
||||
|
||||
[storage]
|
||||
# Storage backend: `filesystem` or `s3`.
|
||||
# If you set `s3` you must also provide `AWS_ACCESS_KEY_ID` and `AWS_SECRET_ACCESS_KEY`.
|
||||
#
|
||||
# Can also be specified via environment variable `BLOB_STORAGE_BACKEND`.
|
||||
#
|
||||
|
||||
Reference in New Issue
Block a user