Skip to main content

backbone_bucket/application/service/
multipart_upload_service.rs

1//! Multipart Upload Service
2//!
3//! Hand-written — NOT generated. This file is safe from regeneration.
4
5use std::sync::Arc;
6
7use chrono::{Duration, Utc};
8use uuid::Uuid;
9
10use super::error::{ServiceError, ServiceResult};
11use crate::domain::entity::{UploadSession, UploadStatus, StorageBackend};
12use crate::infrastructure::persistence::UploadSessionRepository;
13use crate::infrastructure::persistence::BucketRepository;
14use crate::infrastructure::persistence::UserQuotaRepository;
15
16const DEFAULT_SESSION_EXPIRY_HOURS: i64 = 24;
17const DEFAULT_CHUNK_SIZE: i32 = 5 * 1024 * 1024;
18
19pub struct MultipartUploadService {
20    session_repo: Arc<UploadSessionRepository>,
21    bucket_repo: Arc<BucketRepository>,
22    quota_repo: Arc<UserQuotaRepository>,
23}
24
25impl MultipartUploadService {
26    pub fn new(
27        session_repo: Arc<UploadSessionRepository>,
28        bucket_repo: Arc<BucketRepository>,
29        quota_repo: Arc<UserQuotaRepository>,
30    ) -> Self {
31        Self { session_repo, bucket_repo, quota_repo }
32    }
33
34    pub async fn initiate(
35        &self,
36        bucket_id: Uuid,
37        user_id: Uuid,
38        path: &str,
39        filename: &str,
40        mime_type: Option<&str>,
41        file_size: i64,
42        chunk_size: Option<i32>,
43    ) -> ServiceResult<UploadSession> {
44        let _bucket = self.bucket_repo
45            .find_by_id(&bucket_id.to_string())
46            .await
47            .map_err(|e| ServiceError::Repository(backbone_core::RepositoryError::DatabaseError(e.to_string())))?
48            .ok_or(ServiceError::NotFound)?;
49
50        if file_size <= 0 {
51            return Err(ServiceError::Validation("file_size must be positive".into()));
52        }
53
54        let chunk = chunk_size.unwrap_or(DEFAULT_CHUNK_SIZE);
55        if chunk <= 0 {
56            return Err(ServiceError::Validation("chunk_size must be positive".into()));
57        }
58        let total_chunks = ((file_size as f64) / (chunk as f64)).ceil() as i32;
59
60        let mut builder = UploadSession::builder()
61            .bucket_id(bucket_id)
62            .user_id(user_id)
63            .path(path.to_string())
64            .filename(filename.to_string())
65            .file_size(file_size)
66            .chunk_size(chunk)
67            .total_chunks(total_chunks)
68            .uploaded_chunks(0)
69            .status(UploadStatus::Initiated)
70            .storage_backend(StorageBackend::Local)
71            .expires_at(Utc::now() + Duration::hours(DEFAULT_SESSION_EXPIRY_HOURS));
72
73        if let Some(mt) = mime_type {
74            builder = builder.mime_type(mt.to_string());
75        }
76
77        let session = builder.build()
78            .map_err(|e| ServiceError::Validation(e))?;
79
80        let created = self.session_repo
81            .create(&session)
82            .await
83            .map_err(|e| ServiceError::Repository(backbone_core::RepositoryError::DatabaseError(e.to_string())))?;
84
85        Ok(created)
86    }
87
88    pub async fn record_part(
89        &self,
90        session_id: Uuid,
91        user_id: Uuid,
92        part_number: i32,
93    ) -> ServiceResult<UploadSession> {
94        let session = self.get_active_session(session_id, user_id).await?;
95
96        if part_number < 1 || part_number > session.total_chunks {
97            return Err(ServiceError::Validation(
98                format!("part_number must be between 1 and {}", session.total_chunks)
99            ));
100        }
101
102        if session.completed_parts.contains(&part_number) {
103            return Err(ServiceError::AlreadyExists(
104                format!("Part {} already uploaded", part_number)
105            ));
106        }
107
108        // TODO: session_repo.record_part — implement custom repository method
109        // For now just update uploaded_chunks
110        let mut updated_session = session;
111        updated_session.uploaded_chunks += 1;
112        if updated_session.status == UploadStatus::Initiated {
113            updated_session.status = UploadStatus::Uploading;
114        }
115        updated_session.metadata.touch();
116        let id_str = updated_session.id.to_string();
117        self.session_repo.update(&id_str, &updated_session).await.map_err(|e| ServiceError::Repository(backbone_core::RepositoryError::DatabaseError(e.to_string())))?;
118
119        self.session_repo
120            .find_by_id(&session_id.to_string())
121            .await
122            .map_err(|e| ServiceError::Repository(backbone_core::RepositoryError::DatabaseError(e.to_string())))?
123            .ok_or(ServiceError::NotFound)
124    }
125
126    pub async fn complete(
127        &self,
128        session_id: Uuid,
129        user_id: Uuid,
130    ) -> ServiceResult<UploadSession> {
131        let mut session = self.get_active_session(session_id, user_id).await?;
132
133        if session.uploaded_chunks < session.total_chunks {
134            return Err(ServiceError::Validation(format!(
135                "Not all parts uploaded: {}/{} complete",
136                session.uploaded_chunks, session.total_chunks
137            )));
138        }
139
140        // Quota gate. Best-effort: a quota row may not exist for every
141        // user (admin-provisioned, opt-in) — absence means "no limit",
142        // not "deny". When a row exists, hard-reject if completing would
143        // push `used_bytes` past `limit_bytes`. The actual increment
144        // happens in `record_completed_usage` after the bytes land.
145        if let Some(quota) = self
146            .quota_repo
147            .find_by_user_id(user_id)
148            .await
149            .map_err(|e| ServiceError::Repository(backbone_core::RepositoryError::DatabaseError(e.to_string())))?
150        {
151            let projected = quota.used_bytes.saturating_add(session.file_size);
152            if projected > quota.limit_bytes {
153                return Err(ServiceError::Validation(format!(
154                    "user quota exceeded: {} + {} > {} bytes",
155                    quota.used_bytes, session.file_size, quota.limit_bytes
156                )));
157            }
158        }
159
160        session.status = UploadStatus::Completing;
161        session.metadata.touch();
162
163        let id_str = session.id.to_string();
164        self.session_repo
165            .update(&id_str, &session)
166            .await
167            .map_err(|e| ServiceError::Repository(backbone_core::RepositoryError::DatabaseError(e.to_string())))?
168            .ok_or(ServiceError::NotFound)
169    }
170
171    /// Check (without mutating) whether `user_id` has capacity for
172    /// `bytes` more. Returns `Ok(())` when there is no quota row (no
173    /// limit configured) or when `used_bytes + bytes <= limit_bytes`.
174    ///
175    /// Called by the single-shot upload handler before writing to
176    /// storage. The resumable flow calls into [`Self::complete`] which
177    /// performs the same check.
178    pub async fn check_capacity(&self, user_id: Uuid, bytes: i64) -> ServiceResult<()> {
179        let Some(quota) = self
180            .quota_repo
181            .find_by_user_id(user_id)
182            .await
183            .map_err(|e| ServiceError::Repository(backbone_core::RepositoryError::DatabaseError(e.to_string())))?
184        else {
185            return Ok(());
186        };
187        let projected = quota.used_bytes.saturating_add(bytes);
188        if projected > quota.limit_bytes {
189            return Err(ServiceError::Validation(format!(
190                "user quota exceeded: {} + {} > {} bytes",
191                quota.used_bytes, bytes, quota.limit_bytes
192            )));
193        }
194        Ok(())
195    }
196
197    /// Record post-upload usage on the user's quota row. Best-effort:
198    /// when no quota row exists this is a no-op. Call AFTER the storage
199    /// `put` and DB row commit succeed.
200    pub async fn record_completed_usage(
201        &self,
202        user_id: Uuid,
203        bytes: i64,
204    ) -> ServiceResult<()> {
205        let Some(mut quota) = self
206            .quota_repo
207            .find_by_user_id(user_id)
208            .await
209            .map_err(|e| ServiceError::Repository(backbone_core::RepositoryError::DatabaseError(e.to_string())))?
210        else {
211            return Ok(());
212        };
213        quota.used_bytes = quota.used_bytes.saturating_add(bytes);
214        quota.file_count = quota.file_count.saturating_add(1);
215        if quota.used_bytes > quota.peak_usage_bytes {
216            quota.peak_usage_bytes = quota.used_bytes;
217            quota.peak_usage_at = Some(Utc::now());
218        }
219        quota.metadata.touch();
220        let id_str = quota.id.to_string();
221        self.quota_repo
222            .update(&id_str, &quota)
223            .await
224            .map_err(|e| ServiceError::Repository(backbone_core::RepositoryError::DatabaseError(e.to_string())))?;
225        Ok(())
226    }
227
228    pub async fn abort(&self, session_id: Uuid, user_id: Uuid) -> ServiceResult<()> {
229        let mut session = self.get_active_session(session_id, user_id).await?;
230        session.status = UploadStatus::Failed;
231        session.metadata.touch();
232        let id_str = session.id.to_string();
233        self.session_repo.update(&id_str, &session).await.map_err(|e| ServiceError::Repository(backbone_core::RepositoryError::DatabaseError(e.to_string())))?;
234        Ok(())
235    }
236
237    pub async fn mark_completed(&self, session_id: Uuid) -> ServiceResult<()> {
238        let mut session = self.session_repo
239            .find_by_id(&session_id.to_string())
240            .await
241            .map_err(|e| ServiceError::Repository(backbone_core::RepositoryError::DatabaseError(e.to_string())))?
242            .ok_or(ServiceError::NotFound)?;
243
244        session.status = UploadStatus::Completed;
245        session.metadata.touch();
246        let id_str = session.id.to_string();
247        self.session_repo.update(&id_str, &session).await.map_err(|e| ServiceError::Repository(backbone_core::RepositoryError::DatabaseError(e.to_string())))?;
248        Ok(())
249    }
250
251    pub async fn cleanup_expired(&self) -> ServiceResult<u64> {
252        // TODO: session_repo.expire_stale_sessions — implement custom method
253        Err(ServiceError::Internal("expire_stale_sessions not yet implemented".to_string()))
254    }
255
256    pub async fn list_active(&self, _user_id: Uuid, _bucket_id: Uuid) -> ServiceResult<Vec<UploadSession>> {
257        // TODO: session_repo.find_by_user_and_bucket — implement custom method
258        Err(ServiceError::Internal("list_active not yet implemented".to_string()))
259    }
260
261    async fn get_active_session(&self, session_id: Uuid, user_id: Uuid) -> ServiceResult<UploadSession> {
262        // TODO: session_repo.find_active_by_id — implement custom method
263        let session = self.session_repo
264            .find_by_id(&session_id.to_string())
265            .await
266            .map_err(|e| ServiceError::Repository(backbone_core::RepositoryError::DatabaseError(e.to_string())))?
267            .ok_or(ServiceError::NotFound)?;
268
269        if session.user_id != user_id {
270            return Err(ServiceError::Validation("Upload session belongs to a different user".into()));
271        }
272
273        if session.expires_at < Utc::now() {
274            return Err(ServiceError::Validation("Upload session has expired".into()));
275        }
276
277        Ok(session)
278    }
279}