backbone_bucket/presentation/http/
serving.rs1use std::sync::Arc;
21
22use axum::extract::{Extension, Path};
23use axum::http::{header, StatusCode};
24use axum::response::{IntoResponse, Redirect, Response};
25use axum::{routing::get, Json, Router};
26use serde_json::json;
27
28use crate::auth::{ArcAuthzPolicy, AuthExtractor, HasOwnerId};
29use crate::config::{BucketConfig, ServingMode};
30use crate::domain::entity::StoredFile;
31use crate::error::{BucketError, BucketResult};
32use crate::infrastructure::persistence::StoredFileRepository;
33use crate::storage::ObjectStorage;
34
35pub struct ServingContext<Identity>
42where
43 Identity: Send + Sync + 'static,
44{
45 pub storage: Arc<dyn ObjectStorage>,
46 pub file_repo: Arc<StoredFileRepository>,
47 pub authz: ArcAuthzPolicy<Identity>,
48 pub config: Arc<BucketConfig>,
49}
50
51impl<I> Clone for ServingContext<I>
52where
53 I: Send + Sync + 'static,
54{
55 fn clone(&self) -> Self {
56 Self {
57 storage: self.storage.clone(),
58 file_repo: self.file_repo.clone(),
59 authz: self.authz.clone(),
60 config: self.config.clone(),
61 }
62 }
63}
64
65pub fn serving_router<A>(ctx: ServingContext<A>) -> Router
79where
80 A: AuthExtractor<()> + HasOwnerId + 'static,
81 A::Rejection: IntoResponse,
82{
83 Router::new()
84 .route("/*key", get(serve_file::<A>))
85 .layer(Extension(Arc::new(ctx)))
86}
87
88async fn serve_file<A>(
89 Extension(ctx): Extension<Arc<ServingContext<A>>>,
90 identity: A,
91 Path(key): Path<String>,
92) -> Result<Response, BucketError>
93where
94 A: AuthExtractor<()> + HasOwnerId + 'static,
95 A::Rejection: IntoResponse,
96{
97 let file = lookup_by_key(&ctx.file_repo, &key).await?;
98 ctx.authz.ensure_can_read(&identity, &file).await?;
99 build_mode_response(&*ctx.storage, &ctx.config, &file).await
100}
101
102pub(crate) async fn build_mode_response(
108 storage: &dyn ObjectStorage,
109 config: &BucketConfig,
110 file: &StoredFile,
111) -> Result<Response, BucketError> {
112 let serving = &config.serving;
113 match serving.default_mode {
114 ServingMode::Redirect => {
115 let url = storage
116 .presigned_get(&file.storage_key, serving.presigned_ttl)
117 .await?;
118 Ok(Redirect::temporary(url.as_str()).into_response())
119 }
120 ServingMode::Stream => {
121 let bytes = storage.get(&file.storage_key).await?;
122 Ok((
123 StatusCode::OK,
124 [(header::CONTENT_TYPE, file.mime_type.clone())],
125 bytes,
126 )
127 .into_response())
128 }
129 ServingMode::SignedUrl => {
130 let url = storage
131 .presigned_get(&file.storage_key, serving.presigned_ttl)
132 .await?;
133 Ok(Json(json!({ "url": url.to_string() })).into_response())
134 }
135 }
136}
137
138pub async fn lookup_by_key(
143 repo: &StoredFileRepository,
144 key: &str,
145) -> BucketResult<StoredFile> {
146 repo.find_by_text_field("storage_key", key)
147 .await
148 .map_err(|e| BucketError::Other(e.to_string()))?
149 .ok_or(BucketError::NotFound)
150}
151
152#[cfg(test)]
153mod tests {
154 use super::*;
155 use std::time::Duration;
156
157 use async_trait::async_trait;
158 use axum::body::to_bytes;
159 use axum::http::StatusCode;
160 use bytes::Bytes;
161 use url::Url;
162 use uuid::Uuid;
163
164 use crate::config::{BucketConfig, ServingConfig, ServingMode, StorageConfig};
165 use crate::domain::entity::{AuditMetadata, FileStatus, StoredFile};
166 use crate::storage::{ObjectMeta, ObjectStorage};
167
168 struct StubStorage {
173 body: Bytes,
174 presigned: Url,
175 }
176
177 #[async_trait]
178 impl ObjectStorage for StubStorage {
179 async fn put(&self, _key: &str, _body: Bytes, _ct: &str) -> BucketResult<()> {
180 Ok(())
181 }
182 async fn get(&self, _key: &str) -> BucketResult<Bytes> {
183 Ok(self.body.clone())
184 }
185 async fn delete(&self, _key: &str) -> BucketResult<()> {
186 Ok(())
187 }
188 async fn head(&self, key: &str) -> BucketResult<ObjectMeta> {
189 Ok(ObjectMeta {
190 key: key.to_string(),
191 size: self.body.len() as u64,
192 content_type: None,
193 etag: None,
194 last_modified: None,
195 })
196 }
197 async fn presigned_get(&self, _key: &str, _ttl: Duration) -> BucketResult<Url> {
198 Ok(self.presigned.clone())
199 }
200 async fn presigned_put(
201 &self,
202 _key: &str,
203 _ttl: Duration,
204 _ct: &str,
205 ) -> BucketResult<Url> {
206 Ok(self.presigned.clone())
207 }
208 fn public_url(&self, _key: &str) -> Option<Url> {
209 None
210 }
211 }
212
213 fn fixture_file() -> StoredFile {
214 StoredFile::new(
215 Uuid::new_v4(),
216 Uuid::new_v4(),
217 "/docs/a.txt".into(),
218 "a.txt".into(),
219 5,
220 "text/plain".into(),
221 false,
222 false,
223 false,
224 false,
225 false,
226 0,
227 FileStatus::Active,
228 "docs/a.txt".into(),
229 1,
230 0,
231 )
232 }
233
234 fn fixture_config(mode: ServingMode) -> BucketConfig {
235 BucketConfig {
236 enabled: true,
237 storage: StorageConfig::Local {
238 root: "/tmp/bucket-test".into(),
239 base_url: Url::parse("http://localhost/cdn/").unwrap(),
240 signing_secret_env: "UNUSED".into(),
241 },
242 serving: ServingConfig {
243 default_mode: mode,
244 public_prefix: "public/".into(),
245 presigned_ttl: Duration::from_secs(60),
246 },
247 }
248 }
249
250 fn stub() -> StubStorage {
251 StubStorage {
252 body: Bytes::from_static(b"hello"),
253 presigned: Url::parse("https://minio.example/docs/a.txt?X-Amz-Signature=abc").unwrap(),
254 }
255 }
256
257 #[tokio::test]
258 async fn redirect_mode_returns_302_to_presigned_url() {
259 let storage = stub();
260 let cfg = fixture_config(ServingMode::Redirect);
261 let file = fixture_file();
262
263 let resp = build_mode_response(&storage, &cfg, &file).await.unwrap();
264 assert_eq!(resp.status(), StatusCode::TEMPORARY_REDIRECT);
265 let loc = resp.headers().get(header::LOCATION).unwrap().to_str().unwrap();
266 assert_eq!(loc, storage.presigned.as_str());
267 }
268
269 #[tokio::test]
270 async fn stream_mode_returns_body_with_content_type() {
271 let storage = stub();
272 let cfg = fixture_config(ServingMode::Stream);
273 let file = fixture_file();
274
275 let resp = build_mode_response(&storage, &cfg, &file).await.unwrap();
276 assert_eq!(resp.status(), StatusCode::OK);
277 assert_eq!(
278 resp.headers().get(header::CONTENT_TYPE).unwrap(),
279 "text/plain"
280 );
281 let body = to_bytes(resp.into_body(), 1024).await.unwrap();
282 assert_eq!(&body[..], b"hello");
283 }
284
285 #[tokio::test]
286 async fn signed_url_mode_returns_json_with_url() {
287 let storage = stub();
288 let cfg = fixture_config(ServingMode::SignedUrl);
289 let file = fixture_file();
290
291 let resp = build_mode_response(&storage, &cfg, &file).await.unwrap();
292 assert_eq!(resp.status(), StatusCode::OK);
293 let body = to_bytes(resp.into_body(), 4096).await.unwrap();
294 let json: serde_json::Value = serde_json::from_slice(&body).unwrap();
295 assert_eq!(json["url"], storage.presigned.as_str());
296 }
297
298 #[test]
299 fn error_response_masks_internal_details() {
300 let resp = BucketError::S3("AWS secret leak: AKIA...".into()).into_response();
302 assert_eq!(resp.status(), StatusCode::INTERNAL_SERVER_ERROR);
303 }
304
305 #[test]
306 fn error_response_preserves_user_facing_status() {
307 assert_eq!(
308 BucketError::NotFound.into_response().status(),
309 StatusCode::NOT_FOUND
310 );
311 assert_eq!(
312 BucketError::Forbidden.into_response().status(),
313 StatusCode::FORBIDDEN
314 );
315 assert_eq!(
316 BucketError::Unauthenticated.into_response().status(),
317 StatusCode::UNAUTHORIZED
318 );
319 assert_eq!(
320 BucketError::InvalidSignature.into_response().status(),
321 StatusCode::FORBIDDEN
322 );
323 }
324}
325
326impl IntoResponse for BucketError {
327 fn into_response(self) -> Response {
328 match &self {
332 BucketError::NotFound => (StatusCode::NOT_FOUND, "not found".to_string()).into_response(),
333 BucketError::Forbidden => (StatusCode::FORBIDDEN, "forbidden".to_string()).into_response(),
334 BucketError::Unauthenticated => {
335 (StatusCode::UNAUTHORIZED, "unauthenticated".to_string()).into_response()
336 }
337 BucketError::InvalidSignature => {
338 (StatusCode::FORBIDDEN, "invalid signature".to_string()).into_response()
339 }
340 BucketError::Unsupported(msg) => (
341 StatusCode::NOT_IMPLEMENTED,
342 format!("operation not supported: {msg}"),
343 )
344 .into_response(),
345 BucketError::InvalidInput(msg) => {
346 (StatusCode::BAD_REQUEST, msg.clone()).into_response()
347 }
348 BucketError::Conflict(msg) => {
349 (StatusCode::CONFLICT, msg.clone()).into_response()
350 }
351 BucketError::PayloadTooLarge(msg) => {
352 (StatusCode::PAYLOAD_TOO_LARGE, msg.clone()).into_response()
353 }
354 BucketError::UnsupportedMediaType(msg) => {
355 (StatusCode::UNSUPPORTED_MEDIA_TYPE, msg.clone()).into_response()
356 }
357 BucketError::Config(_)
358 | BucketError::Io(_)
359 | BucketError::S3(_)
360 | BucketError::Url(_)
361 | BucketError::Other(_) => {
362 tracing::error!(error = %self, "bucket handler failed");
363 (StatusCode::INTERNAL_SERVER_ERROR, "internal error".to_string())
364 .into_response()
365 }
366 }
367 }
368}