backbone_bucket/storage/
s3.rs1use 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
35pub struct S3Storage {
37 client: Client,
38 cfg: S3Config,
39 serving: ServingConfig,
40}
41
42impl S3Storage {
43 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 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 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 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}