Skip to main content

backbone_bucket/presentation/http/
serving.rs

1//! Mode-B serving handler.
2//!
3//! Serves files under `/*key` with the consumer's auth extractor and
4//! authorization policy in the path. Response strategy is driven by
5//! [`ServingMode`](crate::config::ServingMode):
6//!
7//! - `Redirect` (default) — 302 to a short-lived presigned URL.
8//! - `Stream` — proxy the bytes through the service.
9//! - `SignedUrl` — return `{"url": "..."}` JSON.
10//!
11//! Authentication and authorization are pluggable:
12//!
13//! - [`AuthExtractor`] resolves the caller's identity from the request.
14//! - [`AuthzPolicy`] decides whether that identity may read the file.
15//!
16//! The context object travels as an Axum [`Extension`] layer, which keeps
17//! the Router generic over `()` so consumers can merge the serving router
18//! into their existing app without wrestling with shared-state generics.
19
20use 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
35/// Shared handler state — travels as an Axum `Extension`.
36///
37/// Parameterized over the consumer's identity type (`A`, which in Axum
38/// doubles as both extractor and identity) so the handler can call the
39/// consumer's [`AuthzPolicy`] without the module knowing what the
40/// identity actually looks like.
41pub 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
65/// Build the serving router.
66///
67/// Consumers mount this at whatever path suits their URL layout —
68/// `/cdn` is the default recommendation:
69///
70/// ```ignore
71/// let app = Router::new()
72///     .merge(bucket.crud_router())
73///     .nest("/cdn", bucket.serving_router::<MyAuth>(policy)?);
74/// ```
75///
76/// `A` is the consumer's [`AuthExtractor`]; Axum resolves it on every
77/// request to produce the identity that flows into the authz policy.
78pub 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
102/// Build the serving response for a single file according to the
103/// configured [`ServingMode`].
104///
105/// Split out of `serve_file` so unit tests can exercise each mode
106/// without spinning up a database or Axum router.
107pub(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
138/// Look up a file by its `storage_key` column.
139///
140/// Exposed as a free fn so call sites outside the handler (tests, custom
141/// routes) can reuse the same lookup rule.
142pub 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    /// In-memory ObjectStorage stub for mode-branch tests.
169    ///
170    /// Records the keys it was asked for so assertions can confirm the
171    /// handler called the right method. No filesystem or network.
172    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        // S3 / Other / Io / Url / Config all collapse to a generic 500.
301        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        // 4xx variants carry caller-facing detail; 5xx variants are
329        // collapsed to a generic message to avoid leaking backend
330        // internals. The full error is captured in `tracing` for operators.
331        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}