Skip to main content

backbone_bucket/storage/
s3.rs

1//! S3 / MinIO backend using `aws-sdk-s3`.
2//!
3//! Works against any S3-compatible service (AWS, MinIO, Wasabi, …). For
4//! MinIO, set `force_path_style = true` and point `endpoint` at the MinIO
5//! URL. Credentials are sourced from the env var *names* supplied via
6//! [`S3Config`] — the module never holds the secret at rest, only the
7//! env-var key used to read it at request time.
8//!
9//! # SigV4
10//!
11//! Presigning is delegated to `aws-sdk-s3`'s `presigned_config` /
12//! `presigned` builders, which implement the AWS SigV4 spec in full. This
13//! replaces the ad-hoc HMAC scheme that previously shipped in
14//! [`crate::CdnService`].
15
16use std::time::Duration;
17
18use async_trait::async_trait;
19use aws_config::{BehaviorVersion, Region};
20use aws_credential_types::Credentials;
21use aws_sdk_s3::config::SharedCredentialsProvider;
22use aws_sdk_s3::operation::get_object::GetObjectError;
23use aws_sdk_s3::operation::head_object::HeadObjectError;
24use aws_sdk_s3::presigning::PresigningConfig;
25use aws_sdk_s3::primitives::ByteStream;
26use aws_sdk_s3::Client;
27use bytes::Bytes;
28use chrono::{DateTime, Utc};
29use url::Url;
30
31use crate::config::{S3Config, ServingConfig};
32use crate::error::{BucketError, BucketResult};
33use crate::storage::{ObjectMeta, ObjectStorage};
34
35/// AWS S3 / MinIO backend.
36pub struct S3Storage {
37    client: Client,
38    cfg: S3Config,
39    serving: ServingConfig,
40}
41
42impl S3Storage {
43    /// Build an S3 client from the supplied config.
44    ///
45    /// Reads access/secret keys from the env-var names in `S3Config`.
46    pub fn new(cfg: S3Config, serving: ServingConfig) -> BucketResult<Self> {
47        let access_key = std::env::var(&cfg.access_key_env)
48            .map_err(|_| BucketError::Config(format!("{} not set", cfg.access_key_env)))?;
49        let secret_key = std::env::var(&cfg.secret_key_env)
50            .map_err(|_| BucketError::Config(format!("{} not set", cfg.secret_key_env)))?;
51
52        let creds = Credentials::new(access_key, secret_key, None, None, "backbone-bucket");
53        let s3_conf = aws_sdk_s3::Config::builder()
54            .behavior_version(BehaviorVersion::latest())
55            .region(Region::new(cfg.region.clone()))
56            .endpoint_url(cfg.endpoint.as_str())
57            .credentials_provider(SharedCredentialsProvider::new(creds))
58            .force_path_style(cfg.force_path_style)
59            .build();
60
61        Ok(Self {
62            client: Client::from_conf(s3_conf),
63            cfg,
64            serving,
65        })
66    }
67
68    fn bucket_for(&self, key: &str) -> &str {
69        if self.is_public_key(key) {
70            self.cfg
71                .public_bucket
72                .as_deref()
73                .unwrap_or(&self.cfg.private_bucket)
74        } else {
75            &self.cfg.private_bucket
76        }
77    }
78
79    fn is_public_key(&self, key: &str) -> bool {
80        !self.serving.public_prefix.is_empty()
81            && key.starts_with(&self.serving.public_prefix)
82            && self.cfg.public_bucket.is_some()
83    }
84
85    fn presigning_config(ttl: Duration) -> BucketResult<PresigningConfig> {
86        PresigningConfig::expires_in(ttl)
87            .map_err(|e| BucketError::S3(format!("presigning config: {e}")))
88    }
89}
90
91#[async_trait]
92impl ObjectStorage for S3Storage {
93    async fn put(&self, key: &str, body: Bytes, content_type: &str) -> BucketResult<()> {
94        self.client
95            .put_object()
96            .bucket(self.bucket_for(key))
97            .key(key)
98            .content_type(content_type)
99            .body(ByteStream::from(body))
100            .send()
101            .await
102            .map_err(|e| BucketError::S3(e.to_string()))?;
103        Ok(())
104    }
105
106    async fn get(&self, key: &str) -> BucketResult<Bytes> {
107        let resp = self
108            .client
109            .get_object()
110            .bucket(self.bucket_for(key))
111            .key(key)
112            .send()
113            .await
114            .map_err(|e| match e.into_service_error() {
115                GetObjectError::NoSuchKey(_) => BucketError::NotFound,
116                other => BucketError::S3(other.to_string()),
117            })?;
118        let data = resp
119            .body
120            .collect()
121            .await
122            .map_err(|e| BucketError::S3(e.to_string()))?;
123        Ok(data.into_bytes())
124    }
125
126    async fn delete(&self, key: &str) -> BucketResult<()> {
127        self.client
128            .delete_object()
129            .bucket(self.bucket_for(key))
130            .key(key)
131            .send()
132            .await
133            .map_err(|e| BucketError::S3(e.to_string()))?;
134        Ok(())
135    }
136
137    async fn head(&self, key: &str) -> BucketResult<ObjectMeta> {
138        let resp = self
139            .client
140            .head_object()
141            .bucket(self.bucket_for(key))
142            .key(key)
143            .send()
144            .await
145            .map_err(|e| match e.into_service_error() {
146                HeadObjectError::NotFound(_) => BucketError::NotFound,
147                other => BucketError::S3(other.to_string()),
148            })?;
149
150        let size = resp
151            .content_length()
152            .and_then(|n| u64::try_from(n).ok())
153            .unwrap_or(0);
154        Ok(ObjectMeta {
155            key: key.to_string(),
156            size,
157            content_type: resp.content_type().map(str::to_string),
158            etag: resp.e_tag().map(str::to_string),
159            last_modified: resp
160                .last_modified()
161                .and_then(|t| {
162                    let secs = t.secs();
163                    DateTime::<Utc>::from_timestamp(secs, t.subsec_nanos())
164                }),
165        })
166    }
167
168    async fn presigned_get(&self, key: &str, ttl: Duration) -> BucketResult<Url> {
169        let presigning = Self::presigning_config(ttl)?;
170        let presigned = self
171            .client
172            .get_object()
173            .bucket(self.bucket_for(key))
174            .key(key)
175            .presigned(presigning)
176            .await
177            .map_err(|e| BucketError::S3(format!("presign get: {e}")))?;
178        Url::parse(presigned.uri()).map_err(Into::into)
179    }
180
181    async fn presigned_put(
182        &self,
183        key: &str,
184        ttl: Duration,
185        content_type: &str,
186    ) -> BucketResult<Url> {
187        let presigning = Self::presigning_config(ttl)?;
188        let presigned = self
189            .client
190            .put_object()
191            .bucket(self.bucket_for(key))
192            .key(key)
193            .content_type(content_type)
194            .presigned(presigning)
195            .await
196            .map_err(|e| BucketError::S3(format!("presign put: {e}")))?;
197        Url::parse(presigned.uri()).map_err(Into::into)
198    }
199
200    fn public_url(&self, key: &str) -> Option<Url> {
201        if !self.is_public_key(key) {
202            return None;
203        }
204        let bucket = self.cfg.public_bucket.as_deref()?;
205        let base = self
206            .cfg
207            .public_endpoint
208            .clone()
209            .unwrap_or_else(|| self.cfg.endpoint.clone());
210        if self.cfg.force_path_style {
211            base.join(&format!("{}/{}", bucket, key)).ok()
212        } else {
213            // virtual-hosted-style: bucket becomes a subdomain host.
214            let mut u = base;
215            if let Some(host) = u.host_str() {
216                let new_host = format!("{bucket}.{host}");
217                u.set_host(Some(&new_host)).ok()?;
218            }
219            u.join(key).ok()
220        }
221    }
222}
223
224#[cfg(test)]
225mod tests {
226    use super::*;
227
228    fn test_config() -> (S3Config, ServingConfig) {
229        // SAFETY: test-only. No live S3 calls are made — we only invoke
230        // the synchronous presigner, which reads these via env lookup.
231        std::env::set_var("_TEST_AK", "AKIAIOSFODNN7EXAMPLE");
232        std::env::set_var(
233            "_TEST_SK",
234            "wJalrXUtnFEMI/K7MDENG/bPxRfiCYEXAMPLEKEY",
235        );
236        let cfg = S3Config {
237            endpoint: Url::parse("http://minio:9000").unwrap(),
238            region: "us-east-1".into(),
239            access_key_env: "_TEST_AK".into(),
240            secret_key_env: "_TEST_SK".into(),
241            private_bucket: "private".into(),
242            public_bucket: Some("public".into()),
243            public_endpoint: Some(Url::parse("https://bucket.example.com").unwrap()),
244            force_path_style: true,
245        };
246        let serving = ServingConfig {
247            default_mode: crate::config::ServingMode::Redirect,
248            public_prefix: "public/".into(),
249            presigned_ttl: Duration::from_secs(60),
250        };
251        (cfg, serving)
252    }
253
254    #[tokio::test]
255    async fn presign_url_contains_sigv4_query_params() {
256        let (cfg, serving) = test_config();
257        let s3 = S3Storage::new(cfg, serving).unwrap();
258
259        let url = s3
260            .presigned_get("foo/bar.txt", Duration::from_secs(60))
261            .await
262            .expect("presign should succeed with env creds set");
263
264        let query: std::collections::HashMap<_, _> =
265            url.query_pairs().into_owned().collect();
266
267        // SigV4 canonical query parameters.
268        assert!(query.contains_key("X-Amz-Algorithm"));
269        assert_eq!(query["X-Amz-Algorithm"], "AWS4-HMAC-SHA256");
270        assert!(query.contains_key("X-Amz-Credential"));
271        assert!(query.contains_key("X-Amz-Date"));
272        assert!(query.contains_key("X-Amz-Expires"));
273        assert!(query.contains_key("X-Amz-Signature"));
274        assert_eq!(query["X-Amz-Expires"], "60");
275    }
276
277    #[tokio::test]
278    async fn public_url_routes_public_prefix_only() {
279        let (cfg, serving) = test_config();
280        let s3 = S3Storage::new(cfg, serving).unwrap();
281
282        assert!(s3.public_url("public/foo.jpg").is_some());
283        assert!(s3.public_url("private/foo.jpg").is_none());
284        assert!(s3.public_url("foo.jpg").is_none());
285    }
286}