1use std::collections::HashMap;
8use std::sync::Arc;
9
10use axum::Router;
11use serde::{Deserialize, Serialize};
12use uuid::Uuid;
13use chrono::{DateTime, Utc};
14
15use backbone_core::http::{ApiResponse, BackboneCrudHandler};
17
18#[cfg(feature = "auth")]
20use backbone_auth::middleware::AuthContext;
21#[cfg(feature = "auth")]
22use backbone_auth::AuthMiddleware;
23
24use crate::domain::entity::*;
26use crate::application::service::{UploadSessionService, ServiceError};
27
28use crate::presentation::dto::{CreateUploadSessionDto, UpdateUploadSessionDto, PatchUploadSessionDto, UploadSessionResponseDto};
30
31use crate::domain::state_machine::{UploadSessionState, UploadSessionStateMachine, UploadSessionTransition};
32
33#[derive(Debug, thiserror::Error)]
35pub enum UploadSessionError {
36 #[error("Not found: {0}")]
37 NotFound(String),
38 #[error("Validation error: {0}")]
39 Validation(String),
40 #[error("Database error: {0}")]
41 Database(String),
42 #[error("Internal error: {0}")]
43 Internal(String),
44 #[error("Target bucket must exist: {0}")]
46 BucketNotFound(String),
47 #[error("Bucket must accept uploads: {0}")]
48 BucketNotActive(String),
49 #[error("File size must be positive: {0}")]
50 InvalidFileSize(String),
51 #[error("Chunk size must be at least 5MB: {0}")]
52 ChunkTooSmall(String),
53 #[error("Chunk size must not exceed 5GB: {0}")]
54 ChunkTooLarge(String),
55 #[error("Total chunks must match file size and chunk size: {0}")]
56 InvalidChunkCount(String),
57 #[error("Path must be valid and not contain traversal: {0}")]
58 InvalidPath(String),
59 #[error("Expiration must be in future: {0}")]
60 InvalidExpiration(String),
61 #[error("Expiration must be within 7 days: {0}")]
62 ExpirationTooFar(String),
63 #[error("Session has expired: {0}")]
64 SessionExpired(String),
65}
66
67impl From<ServiceError> for UploadSessionError {
68 fn from(err: ServiceError) -> Self {
69 match err {
70 ServiceError::NotFound => Self::NotFound(err.to_string()),
71 ServiceError::Validation(ref msg) => Self::Validation(msg.clone()),
72 ServiceError::AlreadyExists(ref msg) => Self::Validation(msg.clone()),
73 ServiceError::Repository(ref e) => Self::Database(e.to_string()),
74 ServiceError::Internal(ref msg) => Self::Internal(msg.clone()),
75 ServiceError::Violations(_) => Self::Validation(err.to_string()),
76 }
77 }
78}
79
80impl axum::response::IntoResponse for UploadSessionError {
81 fn into_response(self) -> axum::response::Response {
82 use axum::http::StatusCode;
83 use axum::Json;
84
85 let (status, code) = match &self {
86 Self::NotFound(_) => (StatusCode::NOT_FOUND, "UPLOADSESSION_NOT_FOUND"),
87 Self::Validation(_) => (StatusCode::BAD_REQUEST, "UPLOADSESSION_VALIDATION_ERROR"),
88 Self::Database(_) => (StatusCode::INTERNAL_SERVER_ERROR, "UPLOADSESSION_DATABASE_ERROR"),
89 Self::Internal(_) => (StatusCode::INTERNAL_SERVER_ERROR, "UPLOADSESSION_INTERNAL_ERROR"),
90 Self::BucketNotFound(_) => (StatusCode::UNPROCESSABLE_ENTITY, "UPLOADSESSION_BUCKET_NOT_FOUND"),
91 Self::BucketNotActive(_) => (StatusCode::UNPROCESSABLE_ENTITY, "UPLOADSESSION_BUCKET_NOT_ACTIVE"),
92 Self::InvalidFileSize(_) => (StatusCode::UNPROCESSABLE_ENTITY, "UPLOADSESSION_INVALID_FILE_SIZE"),
93 Self::ChunkTooSmall(_) => (StatusCode::UNPROCESSABLE_ENTITY, "UPLOADSESSION_CHUNK_TOO_SMALL"),
94 Self::ChunkTooLarge(_) => (StatusCode::UNPROCESSABLE_ENTITY, "UPLOADSESSION_CHUNK_TOO_LARGE"),
95 Self::InvalidChunkCount(_) => (StatusCode::UNPROCESSABLE_ENTITY, "UPLOADSESSION_INVALID_CHUNK_COUNT"),
96 Self::InvalidPath(_) => (StatusCode::UNPROCESSABLE_ENTITY, "UPLOADSESSION_INVALID_PATH"),
97 Self::InvalidExpiration(_) => (StatusCode::UNPROCESSABLE_ENTITY, "UPLOADSESSION_INVALID_EXPIRATION"),
98 Self::ExpirationTooFar(_) => (StatusCode::UNPROCESSABLE_ENTITY, "UPLOADSESSION_EXPIRATION_TOO_FAR"),
99 Self::SessionExpired(_) => (StatusCode::UNPROCESSABLE_ENTITY, "UPLOADSESSION_SESSION_EXPIRED"),
100 };
101
102 let body = serde_json::json!({
103 "success": false,
104 "error": code,
105 "message": self.to_string(),
106 });
107
108 (status, Json(body)).into_response()
109 }
110}
111
112pub mod upload_session_errors {
114 pub const BUCKET_NOT_FOUND: &str = "UPLOADSESSION_BUCKET_NOT_FOUND";
115 pub const BUCKET_NOT_ACTIVE: &str = "UPLOADSESSION_BUCKET_NOT_ACTIVE";
116 pub const INVALID_FILE_SIZE: &str = "UPLOADSESSION_INVALID_FILE_SIZE";
117 pub const CHUNK_TOO_SMALL: &str = "UPLOADSESSION_CHUNK_TOO_SMALL";
118 pub const CHUNK_TOO_LARGE: &str = "UPLOADSESSION_CHUNK_TOO_LARGE";
119 pub const INVALID_CHUNK_COUNT: &str = "UPLOADSESSION_INVALID_CHUNK_COUNT";
120 pub const INVALID_PATH: &str = "UPLOADSESSION_INVALID_PATH";
121 pub const INVALID_EXPIRATION: &str = "UPLOADSESSION_INVALID_EXPIRATION";
122 pub const EXPIRATION_TOO_FAR: &str = "UPLOADSESSION_EXPIRATION_TOO_FAR";
123 pub const SESSION_EXPIRED: &str = "UPLOADSESSION_SESSION_EXPIRED";
124}
125
126pub fn create_upload_session_routes(service: Arc<UploadSessionService>) -> Router {
159 BackboneCrudHandler::<UploadSessionService, UploadSession, CreateUploadSessionDto, UpdateUploadSessionDto, UploadSessionResponseDto>::routes(
160 service,
161 "/upload_sessions",
162 )
163}
164
165pub fn create_upload_session_read_routes(service: Arc<UploadSessionService>) -> Router {
171 BackboneCrudHandler::<UploadSessionService, UploadSession, CreateUploadSessionDto, UpdateUploadSessionDto, UploadSessionResponseDto>::read_routes(
172 service,
173 "/upload_sessions",
174 )
175}
176
177pub fn create_upload_session_write_routes(service: Arc<UploadSessionService>) -> Router {
189 BackboneCrudHandler::<UploadSessionService, UploadSession, CreateUploadSessionDto, UpdateUploadSessionDto, UploadSessionResponseDto>::write_routes(
190 service,
191 "/upload_sessions",
192 )
193}
194
195#[cfg(feature = "auth")]
201pub fn create_protected_upload_session_routes<A: AuthMiddleware + Send + Sync + 'static>(
202 service: Arc<UploadSessionService>,
203 auth: Arc<A>,
204) -> Router {
205 use axum::middleware;
206 use axum::response::IntoResponse;
207
208 let auth_layer = auth.clone();
209 create_upload_session_routes(service)
210 .layer(middleware::from_fn(move |mut req: axum::extract::Request, next: axum::middleware::Next| {
211 let auth = auth_layer.clone();
212 async move {
213 let token = req.headers()
214 .get(axum::http::header::AUTHORIZATION)
215 .and_then(|h| h.to_str().ok())
216 .and_then(|raw| raw.strip_prefix("Bearer ").or_else(|| raw.strip_prefix("bearer ")))
217 .unwrap_or("");
218 match auth.authenticate(token).await {
219 Ok(ctx) => {
220 req.extensions_mut().insert(ctx);
221 next.run(req).await
222 }
223 Err(_) => {
224 (axum::http::StatusCode::UNAUTHORIZED,
225 axum::Json(serde_json::json!({
226 "success": false,
227 "error": "unauthorized",
228 "message": "Authentication required"
229 }))
230 ).into_response()
231 }
232 }
233 }
234 }))
235}
236
237pub async fn start_upload_transition(
245 axum::extract::State(service): axum::extract::State<Arc<UploadSessionService>>,
246 axum::extract::Path(id): axum::extract::Path<String>,
247 #[cfg(feature = "auth")] axum::Extension(auth): axum::Extension<AuthContext>,
248) -> impl axum::response::IntoResponse {
249 use axum::{http::StatusCode, Json};
250
251 let entity = match service.get_by_id(&id).await {
253 Ok(Some(e)) => e,
254 Ok(None) => return (StatusCode::NOT_FOUND, Json(ApiResponse::<UploadSessionResponseDto>::not_found("UploadSession", &id))),
255 Err(e) => return (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<UploadSessionResponseDto>::error(e.to_string()))),
256 };
257
258 #[cfg(feature = "auth")]
260 {
261 let allowed_roles = UploadSessionTransition::StartUpload.allowed_roles();
262 let has_specific_perm = auth.permissions.iter().any(|p| p == "upload_session:transition:start_upload");
263 let has_update_perm = auth.permissions.iter().any(|p| p == "upload_session:update");
264 let has_role = auth.roles.iter().any(|r| allowed_roles.contains(&r.as_str()));
265 if !has_specific_perm && !has_update_perm && !has_role {
266 return (StatusCode::FORBIDDEN, Json(ApiResponse::<UploadSessionResponseDto>::error("Insufficient permissions for start_upload transition")));
267 }
268 }
269
270 let current_state: UploadSessionState = entity.status.to_string().parse()
272 .unwrap_or(UploadSessionState::default());
273 let sm = UploadSessionStateMachine::from_state(current_state);
274 if !sm.can_transition(UploadSessionTransition::StartUpload) {
275 return (StatusCode::BAD_REQUEST, Json(ApiResponse::<UploadSessionResponseDto>::error("Transition not allowed from current state")));
276 }
277
278 let mut fields: HashMap<String, serde_json::Value> = HashMap::new();
280 fields.insert("status".to_string(), serde_json::Value::String("Uploading".to_string()));
281
282 match service.partial_update(&id, fields).await {
283 Ok(Some(updated)) => {
284 let response: UploadSessionResponseDto = updated.into();
285 (StatusCode::OK, Json(ApiResponse::ok(response)))
286 }
287 Ok(None) => (StatusCode::NOT_FOUND, Json(ApiResponse::<UploadSessionResponseDto>::not_found("UploadSession", &id))),
288 Err(e) => (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<UploadSessionResponseDto>::error(e.to_string()))),
289 }
290}
291
292pub async fn add_part_transition(
296 axum::extract::State(service): axum::extract::State<Arc<UploadSessionService>>,
297 axum::extract::Path(id): axum::extract::Path<String>,
298 #[cfg(feature = "auth")] axum::Extension(auth): axum::Extension<AuthContext>,
299) -> impl axum::response::IntoResponse {
300 use axum::{http::StatusCode, Json};
301
302 let entity = match service.get_by_id(&id).await {
304 Ok(Some(e)) => e,
305 Ok(None) => return (StatusCode::NOT_FOUND, Json(ApiResponse::<UploadSessionResponseDto>::not_found("UploadSession", &id))),
306 Err(e) => return (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<UploadSessionResponseDto>::error(e.to_string()))),
307 };
308
309 #[cfg(feature = "auth")]
311 {
312 let allowed_roles = UploadSessionTransition::AddPart.allowed_roles();
313 let has_specific_perm = auth.permissions.iter().any(|p| p == "upload_session:transition:add_part");
314 let has_update_perm = auth.permissions.iter().any(|p| p == "upload_session:update");
315 let has_role = auth.roles.iter().any(|r| allowed_roles.contains(&r.as_str()));
316 if !has_specific_perm && !has_update_perm && !has_role {
317 return (StatusCode::FORBIDDEN, Json(ApiResponse::<UploadSessionResponseDto>::error("Insufficient permissions for add_part transition")));
318 }
319 }
320
321 let current_state: UploadSessionState = entity.status.to_string().parse()
323 .unwrap_or(UploadSessionState::default());
324 let sm = UploadSessionStateMachine::from_state(current_state);
325 if !sm.can_transition(UploadSessionTransition::AddPart) {
326 return (StatusCode::BAD_REQUEST, Json(ApiResponse::<UploadSessionResponseDto>::error("Transition not allowed from current state")));
327 }
328
329 let mut fields: HashMap<String, serde_json::Value> = HashMap::new();
331 fields.insert("status".to_string(), serde_json::Value::String("Uploading".to_string()));
332
333 match service.partial_update(&id, fields).await {
334 Ok(Some(updated)) => {
335 let response: UploadSessionResponseDto = updated.into();
336 (StatusCode::OK, Json(ApiResponse::ok(response)))
337 }
338 Ok(None) => (StatusCode::NOT_FOUND, Json(ApiResponse::<UploadSessionResponseDto>::not_found("UploadSession", &id))),
339 Err(e) => (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<UploadSessionResponseDto>::error(e.to_string()))),
340 }
341}
342
343pub async fn complete_transition(
347 axum::extract::State(service): axum::extract::State<Arc<UploadSessionService>>,
348 axum::extract::Path(id): axum::extract::Path<String>,
349 #[cfg(feature = "auth")] axum::Extension(auth): axum::Extension<AuthContext>,
350) -> impl axum::response::IntoResponse {
351 use axum::{http::StatusCode, Json};
352
353 let entity = match service.get_by_id(&id).await {
355 Ok(Some(e)) => e,
356 Ok(None) => return (StatusCode::NOT_FOUND, Json(ApiResponse::<UploadSessionResponseDto>::not_found("UploadSession", &id))),
357 Err(e) => return (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<UploadSessionResponseDto>::error(e.to_string()))),
358 };
359
360 #[cfg(feature = "auth")]
362 {
363 let allowed_roles = UploadSessionTransition::Complete.allowed_roles();
364 let has_specific_perm = auth.permissions.iter().any(|p| p == "upload_session:transition:complete");
365 let has_update_perm = auth.permissions.iter().any(|p| p == "upload_session:update");
366 let has_role = auth.roles.iter().any(|r| allowed_roles.contains(&r.as_str()));
367 if !has_specific_perm && !has_update_perm && !has_role {
368 return (StatusCode::FORBIDDEN, Json(ApiResponse::<UploadSessionResponseDto>::error("Insufficient permissions for complete transition")));
369 }
370 }
371
372 let current_state: UploadSessionState = entity.status.to_string().parse()
374 .unwrap_or(UploadSessionState::default());
375 let sm = UploadSessionStateMachine::from_state(current_state);
376 if !sm.can_transition(UploadSessionTransition::Complete) {
377 return (StatusCode::BAD_REQUEST, Json(ApiResponse::<UploadSessionResponseDto>::error("Transition not allowed from current state")));
378 }
379
380 let mut fields: HashMap<String, serde_json::Value> = HashMap::new();
382 fields.insert("status".to_string(), serde_json::Value::String("Completing".to_string()));
383
384 match service.partial_update(&id, fields).await {
385 Ok(Some(updated)) => {
386 let response: UploadSessionResponseDto = updated.into();
387 (StatusCode::OK, Json(ApiResponse::ok(response)))
388 }
389 Ok(None) => (StatusCode::NOT_FOUND, Json(ApiResponse::<UploadSessionResponseDto>::not_found("UploadSession", &id))),
390 Err(e) => (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<UploadSessionResponseDto>::error(e.to_string()))),
391 }
392}
393
394pub async fn finish_transition(
398 axum::extract::State(service): axum::extract::State<Arc<UploadSessionService>>,
399 axum::extract::Path(id): axum::extract::Path<String>,
400 #[cfg(feature = "auth")] axum::Extension(auth): axum::Extension<AuthContext>,
401) -> impl axum::response::IntoResponse {
402 use axum::{http::StatusCode, Json};
403
404 let entity = match service.get_by_id(&id).await {
406 Ok(Some(e)) => e,
407 Ok(None) => return (StatusCode::NOT_FOUND, Json(ApiResponse::<UploadSessionResponseDto>::not_found("UploadSession", &id))),
408 Err(e) => return (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<UploadSessionResponseDto>::error(e.to_string()))),
409 };
410
411 #[cfg(feature = "auth")]
413 {
414 let allowed_roles = UploadSessionTransition::Finish.allowed_roles();
415 let has_specific_perm = auth.permissions.iter().any(|p| p == "upload_session:transition:finish");
416 let has_update_perm = auth.permissions.iter().any(|p| p == "upload_session:update");
417 let has_role = auth.roles.iter().any(|r| allowed_roles.contains(&r.as_str()));
418 if !has_specific_perm && !has_update_perm && !has_role {
419 return (StatusCode::FORBIDDEN, Json(ApiResponse::<UploadSessionResponseDto>::error("Insufficient permissions for finish transition")));
420 }
421 }
422
423 let current_state: UploadSessionState = entity.status.to_string().parse()
425 .unwrap_or(UploadSessionState::default());
426 let sm = UploadSessionStateMachine::from_state(current_state);
427 if !sm.can_transition(UploadSessionTransition::Finish) {
428 return (StatusCode::BAD_REQUEST, Json(ApiResponse::<UploadSessionResponseDto>::error("Transition not allowed from current state")));
429 }
430
431 let mut fields: HashMap<String, serde_json::Value> = HashMap::new();
433 fields.insert("status".to_string(), serde_json::Value::String("Completed".to_string()));
434
435 match service.partial_update(&id, fields).await {
436 Ok(Some(updated)) => {
437 let response: UploadSessionResponseDto = updated.into();
438 (StatusCode::OK, Json(ApiResponse::ok(response)))
439 }
440 Ok(None) => (StatusCode::NOT_FOUND, Json(ApiResponse::<UploadSessionResponseDto>::not_found("UploadSession", &id))),
441 Err(e) => (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<UploadSessionResponseDto>::error(e.to_string()))),
442 }
443}
444
445pub async fn fail_transition(
449 axum::extract::State(service): axum::extract::State<Arc<UploadSessionService>>,
450 axum::extract::Path(id): axum::extract::Path<String>,
451 #[cfg(feature = "auth")] axum::Extension(auth): axum::Extension<AuthContext>,
452) -> impl axum::response::IntoResponse {
453 use axum::{http::StatusCode, Json};
454
455 let entity = match service.get_by_id(&id).await {
457 Ok(Some(e)) => e,
458 Ok(None) => return (StatusCode::NOT_FOUND, Json(ApiResponse::<UploadSessionResponseDto>::not_found("UploadSession", &id))),
459 Err(e) => return (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<UploadSessionResponseDto>::error(e.to_string()))),
460 };
461
462 #[cfg(feature = "auth")]
464 {
465 let allowed_roles = UploadSessionTransition::Fail.allowed_roles();
466 let has_specific_perm = auth.permissions.iter().any(|p| p == "upload_session:transition:fail");
467 let has_update_perm = auth.permissions.iter().any(|p| p == "upload_session:update");
468 let has_role = auth.roles.iter().any(|r| allowed_roles.contains(&r.as_str()));
469 if !has_specific_perm && !has_update_perm && !has_role {
470 return (StatusCode::FORBIDDEN, Json(ApiResponse::<UploadSessionResponseDto>::error("Insufficient permissions for fail transition")));
471 }
472 }
473
474 let current_state: UploadSessionState = entity.status.to_string().parse()
476 .unwrap_or(UploadSessionState::default());
477 let sm = UploadSessionStateMachine::from_state(current_state);
478 if !sm.can_transition(UploadSessionTransition::Fail) {
479 return (StatusCode::BAD_REQUEST, Json(ApiResponse::<UploadSessionResponseDto>::error("Transition not allowed from current state")));
480 }
481
482 let mut fields: HashMap<String, serde_json::Value> = HashMap::new();
484 fields.insert("status".to_string(), serde_json::Value::String("Failed".to_string()));
485
486 match service.partial_update(&id, fields).await {
487 Ok(Some(updated)) => {
488 let response: UploadSessionResponseDto = updated.into();
489 (StatusCode::OK, Json(ApiResponse::ok(response)))
490 }
491 Ok(None) => (StatusCode::NOT_FOUND, Json(ApiResponse::<UploadSessionResponseDto>::not_found("UploadSession", &id))),
492 Err(e) => (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<UploadSessionResponseDto>::error(e.to_string()))),
493 }
494}
495
496pub async fn abort_transition(
500 axum::extract::State(service): axum::extract::State<Arc<UploadSessionService>>,
501 axum::extract::Path(id): axum::extract::Path<String>,
502 #[cfg(feature = "auth")] axum::Extension(auth): axum::Extension<AuthContext>,
503) -> impl axum::response::IntoResponse {
504 use axum::{http::StatusCode, Json};
505
506 let entity = match service.get_by_id(&id).await {
508 Ok(Some(e)) => e,
509 Ok(None) => return (StatusCode::NOT_FOUND, Json(ApiResponse::<UploadSessionResponseDto>::not_found("UploadSession", &id))),
510 Err(e) => return (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<UploadSessionResponseDto>::error(e.to_string()))),
511 };
512
513 #[cfg(feature = "auth")]
515 {
516 let allowed_roles = UploadSessionTransition::Abort.allowed_roles();
517 let has_specific_perm = auth.permissions.iter().any(|p| p == "upload_session:transition:abort");
518 let has_update_perm = auth.permissions.iter().any(|p| p == "upload_session:update");
519 let has_role = auth.roles.iter().any(|r| allowed_roles.contains(&r.as_str()));
520 if !has_specific_perm && !has_update_perm && !has_role {
521 return (StatusCode::FORBIDDEN, Json(ApiResponse::<UploadSessionResponseDto>::error("Insufficient permissions for abort transition")));
522 }
523 }
524
525 let current_state: UploadSessionState = entity.status.to_string().parse()
527 .unwrap_or(UploadSessionState::default());
528 let sm = UploadSessionStateMachine::from_state(current_state);
529 if !sm.can_transition(UploadSessionTransition::Abort) {
530 return (StatusCode::BAD_REQUEST, Json(ApiResponse::<UploadSessionResponseDto>::error("Transition not allowed from current state")));
531 }
532
533 let mut fields: HashMap<String, serde_json::Value> = HashMap::new();
535 fields.insert("status".to_string(), serde_json::Value::String("Failed".to_string()));
536
537 match service.partial_update(&id, fields).await {
538 Ok(Some(updated)) => {
539 let response: UploadSessionResponseDto = updated.into();
540 (StatusCode::OK, Json(ApiResponse::ok(response)))
541 }
542 Ok(None) => (StatusCode::NOT_FOUND, Json(ApiResponse::<UploadSessionResponseDto>::not_found("UploadSession", &id))),
543 Err(e) => (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<UploadSessionResponseDto>::error(e.to_string()))),
544 }
545}
546
547pub async fn expire_transition(
551 axum::extract::State(service): axum::extract::State<Arc<UploadSessionService>>,
552 axum::extract::Path(id): axum::extract::Path<String>,
553 #[cfg(feature = "auth")] axum::Extension(auth): axum::Extension<AuthContext>,
554) -> impl axum::response::IntoResponse {
555 use axum::{http::StatusCode, Json};
556
557 let entity = match service.get_by_id(&id).await {
559 Ok(Some(e)) => e,
560 Ok(None) => return (StatusCode::NOT_FOUND, Json(ApiResponse::<UploadSessionResponseDto>::not_found("UploadSession", &id))),
561 Err(e) => return (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<UploadSessionResponseDto>::error(e.to_string()))),
562 };
563
564 #[cfg(feature = "auth")]
566 {
567 let allowed_roles = UploadSessionTransition::Expire.allowed_roles();
568 let has_specific_perm = auth.permissions.iter().any(|p| p == "upload_session:transition:expire");
569 let has_update_perm = auth.permissions.iter().any(|p| p == "upload_session:update");
570 let has_role = auth.roles.iter().any(|r| allowed_roles.contains(&r.as_str()));
571 if !has_specific_perm && !has_update_perm && !has_role {
572 return (StatusCode::FORBIDDEN, Json(ApiResponse::<UploadSessionResponseDto>::error("Insufficient permissions for expire transition")));
573 }
574 }
575
576 let current_state: UploadSessionState = entity.status.to_string().parse()
578 .unwrap_or(UploadSessionState::default());
579 let sm = UploadSessionStateMachine::from_state(current_state);
580 if !sm.can_transition(UploadSessionTransition::Expire) {
581 return (StatusCode::BAD_REQUEST, Json(ApiResponse::<UploadSessionResponseDto>::error("Transition not allowed from current state")));
582 }
583
584 let mut fields: HashMap<String, serde_json::Value> = HashMap::new();
586 fields.insert("status".to_string(), serde_json::Value::String("Expired".to_string()));
587
588 match service.partial_update(&id, fields).await {
589 Ok(Some(updated)) => {
590 let response: UploadSessionResponseDto = updated.into();
591 (StatusCode::OK, Json(ApiResponse::ok(response)))
592 }
593 Ok(None) => (StatusCode::NOT_FOUND, Json(ApiResponse::<UploadSessionResponseDto>::not_found("UploadSession", &id))),
594 Err(e) => (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<UploadSessionResponseDto>::error(e.to_string()))),
595 }
596}
597
598pub async fn retry_transition(
602 axum::extract::State(service): axum::extract::State<Arc<UploadSessionService>>,
603 axum::extract::Path(id): axum::extract::Path<String>,
604 #[cfg(feature = "auth")] axum::Extension(auth): axum::Extension<AuthContext>,
605) -> impl axum::response::IntoResponse {
606 use axum::{http::StatusCode, Json};
607
608 let entity = match service.get_by_id(&id).await {
610 Ok(Some(e)) => e,
611 Ok(None) => return (StatusCode::NOT_FOUND, Json(ApiResponse::<UploadSessionResponseDto>::not_found("UploadSession", &id))),
612 Err(e) => return (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<UploadSessionResponseDto>::error(e.to_string()))),
613 };
614
615 #[cfg(feature = "auth")]
617 {
618 let allowed_roles = UploadSessionTransition::Retry.allowed_roles();
619 let has_specific_perm = auth.permissions.iter().any(|p| p == "upload_session:transition:retry");
620 let has_update_perm = auth.permissions.iter().any(|p| p == "upload_session:update");
621 let has_role = auth.roles.iter().any(|r| allowed_roles.contains(&r.as_str()));
622 if !has_specific_perm && !has_update_perm && !has_role {
623 return (StatusCode::FORBIDDEN, Json(ApiResponse::<UploadSessionResponseDto>::error("Insufficient permissions for retry transition")));
624 }
625 }
626
627 let current_state: UploadSessionState = entity.status.to_string().parse()
629 .unwrap_or(UploadSessionState::default());
630 let sm = UploadSessionStateMachine::from_state(current_state);
631 if !sm.can_transition(UploadSessionTransition::Retry) {
632 return (StatusCode::BAD_REQUEST, Json(ApiResponse::<UploadSessionResponseDto>::error("Transition not allowed from current state")));
633 }
634
635 let mut fields: HashMap<String, serde_json::Value> = HashMap::new();
637 fields.insert("status".to_string(), serde_json::Value::String("Initiated".to_string()));
638
639 match service.partial_update(&id, fields).await {
640 Ok(Some(updated)) => {
641 let response: UploadSessionResponseDto = updated.into();
642 (StatusCode::OK, Json(ApiResponse::ok(response)))
643 }
644 Ok(None) => (StatusCode::NOT_FOUND, Json(ApiResponse::<UploadSessionResponseDto>::not_found("UploadSession", &id))),
645 Err(e) => (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<UploadSessionResponseDto>::error(e.to_string()))),
646 }
647}
648
649pub fn create_upload_session_transition_routes(service: Arc<UploadSessionService>) -> Router {
651 use axum::routing::post;
652
653 Router::new()
654 .route("/upload_sessions/:id/transitions/start_upload", post(start_upload_transition))
655 .route("/upload_sessions/:id/transitions/add_part", post(add_part_transition))
656 .route("/upload_sessions/:id/transitions/complete", post(complete_transition))
657 .route("/upload_sessions/:id/transitions/finish", post(finish_transition))
658 .route("/upload_sessions/:id/transitions/fail", post(fail_transition))
659 .route("/upload_sessions/:id/transitions/abort", post(abort_transition))
660 .route("/upload_sessions/:id/transitions/expire", post(expire_transition))
661 .route("/upload_sessions/:id/transitions/retry", post(retry_transition))
662 .with_state(service)
663}