Skip to main content

max / makenotwork

9.9 KB · 273 lines History Blame Raw
1 //! Presigned upload and confirm handlers for version files.
2
3 use axum::{
4 extract::{Path, State},
5 response::IntoResponse,
6 Json,
7 };
8 use serde::Deserialize;
9
10 use crate::{
11 auth::AuthUser,
12 db::{self, VersionId},
13 error::{AppError, Result, ResultExt},
14 storage::{FileType, S3Client, CACHE_CONTROL_IMMUTABLE},
15 AppState,
16 };
17
18 use super::{commit_upload, CommitTarget, ConfirmUploadResponse, PresignUploadResponse};
19
20 /// JSON input for requesting a presigned version upload URL.
21 #[derive(Debug, Deserialize)]
22 pub struct VersionPresignRequest {
23 pub file_name: String,
24 pub content_type: String,
25 }
26
27 /// JSON input for confirming a completed version upload.
28 #[derive(Debug, Deserialize)]
29 pub struct VersionConfirmRequest {
30 pub s3_key: String,
31 }
32
33 /// Generate a presigned URL for uploading a version file to S3
34 ///
35 /// POST /api/versions/{version_id}/upload/presign
36 ///
37 /// Requires authentication. User must own the item (through version -> item -> project chain).
38 #[tracing::instrument(skip_all, name = "storage::version_presign_upload")]
39 pub(super) async fn version_presign_upload(
40 State(state): State<AppState>,
41 AuthUser(user): AuthUser,
42 Path(version_id): Path<VersionId>,
43 Json(req): Json<VersionPresignRequest>,
44 ) -> Result<impl IntoResponse> {
45 user.check_not_suspended()?;
46 let s3 = state.require_s3()?;
47
48 let file_type = FileType::Download;
49
50 // Validate content type and extension
51 S3Client::validate_content_type(file_type, &req.content_type)?;
52 S3Client::validate_extension(file_type, &req.file_name)?;
53
54 // Fetch version and verify ownership through version -> item -> project chain
55 let version = db::versions::get_version_by_id(&state.db, version_id)
56 .await?
57 .ok_or(AppError::NotFound)?;
58
59 let owner = db::items::get_item_owner(&state.db, version.item_id)
60 .await?
61 .ok_or(AppError::NotFound)?;
62
63 if owner != user.id {
64 return Err(AppError::Forbidden);
65 }
66
67 // Early quota check
68 db::creator_tiers::check_presign_allowed(&state.db, user.id, file_type).await?;
69
70 let max_file_bytes = db::creator_tiers::get_effective_max_file_bytes(&state.db, user.id, file_type).await?;
71
72 // Generate S3 key using the version's item_id
73 let s3_key = S3Client::generate_key(user.id, version.item_id, file_type, &req.file_name);
74
75 // Track the pending upload so the reaper can clean it up if never confirmed
76 db::pending_uploads::record_pending_upload(&state.db, user.id, &s3_key, "main").await?;
77
78 let expires_in = 3600;
79 let upload_url = s3.presign_upload(&s3_key, &req.content_type, Some(expires_in), Some(CACHE_CONTROL_IMMUTABLE), None)
80 .await
81 .context("presign upload for version file")?;
82
83 Ok(Json(PresignUploadResponse {
84 upload_url,
85 s3_key,
86 expires_in,
87 cache_control: Some(CACHE_CONTROL_IMMUTABLE.to_string()),
88 max_file_bytes,
89 }))
90 }
91
92 /// Confirm that a version file upload has completed and update the database
93 ///
94 /// POST /api/versions/{version_id}/upload/confirm
95 ///
96 /// Requires authentication. User must own the item.
97 #[tracing::instrument(skip_all, name = "storage::version_confirm_upload")]
98 pub(super) async fn version_confirm_upload(
99 State(state): State<AppState>,
100 AuthUser(user): AuthUser,
101 Path(version_id): Path<VersionId>,
102 Json(req): Json<VersionConfirmRequest>,
103 ) -> Result<impl IntoResponse> {
104 user.check_not_suspended()?;
105 let s3 = state.require_s3()?;
106
107 // Fetch version and verify ownership
108 let version = db::versions::get_version_by_id(&state.db, version_id)
109 .await?
110 .ok_or(AppError::NotFound)?;
111
112 let owner = db::items::get_item_owner(&state.db, version.item_id)
113 .await?
114 .ok_or(AppError::NotFound)?;
115
116 if owner != user.id {
117 return Err(AppError::Forbidden);
118 }
119
120 // Validate S3 key belongs to this user + item (prevent cross-user file reference)
121 let expected_prefix = format!("{}/{}/", user.id, version.item_id);
122 if !req.s3_key.starts_with(&expected_prefix) {
123 return Err(AppError::BadRequest(
124 "Invalid upload key".to_string(),
125 ));
126 }
127
128 // Verify the object exists in S3
129 if !s3.object_exists(&req.s3_key).await? {
130 return Err(AppError::BadRequest(
131 "Upload not found. Please try uploading again.".to_string(),
132 ));
133 }
134
135 // Enforce file size limit (versions are always downloads)
136 let file_size_bytes = s3.object_size(&req.s3_key).await?.ok_or_else(|| {
137 AppError::BadRequest("Could not determine file size. Please try uploading again.".to_string())
138 })?;
139 if file_size_bytes as u64 > FileType::Download.max_size() {
140 s3.delete_object(&req.s3_key).await.ok();
141 let limit_mb = FileType::Download.max_size() / (1024 * 1024);
142 let file_mb = file_size_bytes as u64 / (1024 * 1024);
143 return Err(AppError::FileTooLarge(format!(
144 "File is {} MB but the maximum for download files is {} MB.",
145 file_mb, limit_mb
146 )));
147 }
148
149 // Enforce tier-based limits (per-file + storage cap)
150 let max_storage = match db::creator_tiers::check_upload_allowed(&state.db, user.id, FileType::Download, file_size_bytes).await {
151 Ok(max) => max,
152 Err(e) => {
153 s3.delete_object(&req.s3_key).await.ok();
154 return Err(e);
155 }
156 };
157
158 // Idempotency: if the version already has this exact s3_key, return success (no-op).
159 // Must come BEFORE scan enqueue / scan_status flip — re-confirming an already-Clean
160 // version must not knock it back to Pending.
161 if version.s3_key.as_deref() == Some(&req.s3_key) {
162 // Still clear pending_uploads — orphan reaper would otherwise delete
163 // the live S3 object 24h later (Run #7 HIGH-1).
164 if let Err(e) = db::pending_uploads::remove_pending_upload(&state.db, user.id, &req.s3_key).await {
165 tracing::warn!(error = ?e, key = %req.s3_key, "remove_pending_upload failed on idempotent re-confirm");
166 }
167 return Ok(Json(ConfirmUploadResponse { success: true, pending_review: None }));
168 }
169
170 // Atomically check storage cap and update counter BEFORE writing
171 // the version record. Avoids orphaned unbilled file references.
172 let old_s3_key = version.s3_key.clone();
173 let old_size = version.file_size_bytes.unwrap_or(0);
174 if old_s3_key.is_some() && old_size > 0 {
175 // Replacing an existing file: atomic decrement-old + increment-new in one query
176 if let Err(e) = db::creator_tiers::try_replace_storage(&state.db, user.id, old_size, file_size_bytes, max_storage).await {
177 s3.delete_object(&req.s3_key).await.ok();
178 return Err(e);
179 }
180 } else {
181 // Fresh upload (no previous file)
182 if let Err(e) = db::creator_tiers::try_increment_storage(&state.db, user.id, file_size_bytes, max_storage).await {
183 s3.delete_object(&req.s3_key).await.ok();
184 return Err(e);
185 }
186 }
187
188 // Clear the pending upload record now that the upload is confirmed
189 db::pending_uploads::remove_pending_upload(&state.db, user.id, &req.s3_key).await?;
190
191 // Extract file name from the s3_key (last path segment)
192 let file_name = req.s3_key.rsplit('/').next().map(|s| s.to_string());
193
194 // Update version with S3 key, file size, and file name. The expected-old
195 // guard rejects the UPDATE if another confirm raced ahead, so we can roll
196 // back the storage credit + delete the (now orphaned) new S3 object.
197 let updated = db::versions::update_version_file(
198 &state.db,
199 version_id,
200 old_s3_key.as_deref(),
201 &req.s3_key,
202 Some(file_size_bytes),
203 file_name.as_deref(),
204 )
205 .await;
206
207 let updated = match updated {
208 Ok(Some(v)) => v,
209 Ok(None) => {
210 // Lost race: another confirm wrote a different s3_key first.
211 // Roll back the storage change we made above.
212 if old_s3_key.is_some() && old_size > 0 {
213 db::creator_tiers::try_replace_storage(
214 &state.db, user.id, file_size_bytes, old_size, i64::MAX,
215 ).await.ok();
216 } else {
217 db::creator_tiers::decrement_storage_used(&state.db, user.id, file_size_bytes).await.ok();
218 }
219 s3.delete_object(&req.s3_key).await.ok();
220 return Err(AppError::BadRequest(
221 "Version was modified concurrently. Please try uploading again.".to_string(),
222 ));
223 }
224 Err(e) => {
225 // DB error: same rollback as the lost-race case.
226 if old_s3_key.is_some() && old_size > 0 {
227 db::creator_tiers::try_replace_storage(
228 &state.db, user.id, file_size_bytes, old_size, i64::MAX,
229 ).await.ok();
230 } else {
231 db::creator_tiers::decrement_storage_used(&state.db, user.id, file_size_bytes).await.ok();
232 }
233 s3.delete_object(&req.s3_key).await.ok();
234 return Err(e);
235 }
236 };
237 let _ = updated;
238
239 let scan_status = commit_upload(
240 &state,
241 CommitTarget::Version(version_id),
242 &req.s3_key,
243 FileType::Download,
244 user.id,
245 file_size_bytes,
246 ).await?;
247
248 // Enqueue old S3 key for deletion now that the DB record points to the new key
249 if let Some(old_key) = old_s3_key
250 && let Err(e) = db::pending_s3_deletions::enqueue_deletions(
251 &state.db,
252 &[(old_key, "main".to_string())],
253 "version_replace",
254 ).await
255 {
256 tracing::warn!(error = ?e, "failed to enqueue old version S3 key for deletion");
257 }
258
259 tracing::info!(
260 "Version upload confirmed: version={}, key={}, size={}",
261 version_id,
262 req.s3_key,
263 file_size_bytes
264 );
265
266 let pending_review = if scan_status == db::FileScanStatus::HeldForReview {
267 Some(true)
268 } else {
269 None
270 };
271 Ok(Json(ConfirmUploadResponse { success: true, pending_review }))
272 }
273