use std::collections::HashMap;
use std::sync::Arc;
use axum::Router;
use serde::{Deserialize, Serialize};
use uuid::Uuid;
use chrono::{DateTime, Utc};
use backbone_core::http::{ApiResponse, BackboneCrudHandler};
#[cfg(feature = "auth")]
use backbone_auth::middleware::AuthContext;
#[cfg(feature = "auth")]
use backbone_auth::AuthMiddleware;
use crate::domain::entity::*;
use crate::application::service::{UploadSessionService, ServiceError};
use crate::presentation::dto::{CreateUploadSessionDto, UpdateUploadSessionDto, PatchUploadSessionDto, UploadSessionResponseDto};
use crate::domain::state_machine::{UploadSessionState, UploadSessionStateMachine, UploadSessionTransition};
#[derive(Debug, thiserror::Error)]
pub enum UploadSessionError {
#[error("Not found: {0}")]
NotFound(String),
#[error("Validation error: {0}")]
Validation(String),
#[error("Database error: {0}")]
Database(String),
#[error("Internal error: {0}")]
Internal(String),
#[error("Target bucket must exist: {0}")]
BucketNotFound(String),
#[error("Bucket must accept uploads: {0}")]
BucketNotActive(String),
#[error("File size must be positive: {0}")]
InvalidFileSize(String),
#[error("Chunk size must be at least 5MB: {0}")]
ChunkTooSmall(String),
#[error("Chunk size must not exceed 5GB: {0}")]
ChunkTooLarge(String),
#[error("Total chunks must match file size and chunk size: {0}")]
InvalidChunkCount(String),
#[error("Path must be valid and not contain traversal: {0}")]
InvalidPath(String),
#[error("Expiration must be in future: {0}")]
InvalidExpiration(String),
#[error("Expiration must be within 7 days: {0}")]
ExpirationTooFar(String),
#[error("Session has expired: {0}")]
SessionExpired(String),
}
impl From<ServiceError> for UploadSessionError {
fn from(err: ServiceError) -> Self {
match err {
ServiceError::NotFound => Self::NotFound(err.to_string()),
ServiceError::Validation(ref msg) => Self::Validation(msg.clone()),
ServiceError::AlreadyExists(ref msg) => Self::Validation(msg.clone()),
ServiceError::Repository(ref e) => Self::Database(e.to_string()),
ServiceError::Internal(ref msg) => Self::Internal(msg.clone()),
ServiceError::Violations(_) => Self::Validation(err.to_string()),
}
}
}
impl axum::response::IntoResponse for UploadSessionError {
fn into_response(self) -> axum::response::Response {
use axum::http::StatusCode;
use axum::Json;
let (status, code) = match &self {
Self::NotFound(_) => (StatusCode::NOT_FOUND, "UPLOADSESSION_NOT_FOUND"),
Self::Validation(_) => (StatusCode::BAD_REQUEST, "UPLOADSESSION_VALIDATION_ERROR"),
Self::Database(_) => (StatusCode::INTERNAL_SERVER_ERROR, "UPLOADSESSION_DATABASE_ERROR"),
Self::Internal(_) => (StatusCode::INTERNAL_SERVER_ERROR, "UPLOADSESSION_INTERNAL_ERROR"),
Self::BucketNotFound(_) => (StatusCode::UNPROCESSABLE_ENTITY, "UPLOADSESSION_BUCKET_NOT_FOUND"),
Self::BucketNotActive(_) => (StatusCode::UNPROCESSABLE_ENTITY, "UPLOADSESSION_BUCKET_NOT_ACTIVE"),
Self::InvalidFileSize(_) => (StatusCode::UNPROCESSABLE_ENTITY, "UPLOADSESSION_INVALID_FILE_SIZE"),
Self::ChunkTooSmall(_) => (StatusCode::UNPROCESSABLE_ENTITY, "UPLOADSESSION_CHUNK_TOO_SMALL"),
Self::ChunkTooLarge(_) => (StatusCode::UNPROCESSABLE_ENTITY, "UPLOADSESSION_CHUNK_TOO_LARGE"),
Self::InvalidChunkCount(_) => (StatusCode::UNPROCESSABLE_ENTITY, "UPLOADSESSION_INVALID_CHUNK_COUNT"),
Self::InvalidPath(_) => (StatusCode::UNPROCESSABLE_ENTITY, "UPLOADSESSION_INVALID_PATH"),
Self::InvalidExpiration(_) => (StatusCode::UNPROCESSABLE_ENTITY, "UPLOADSESSION_INVALID_EXPIRATION"),
Self::ExpirationTooFar(_) => (StatusCode::UNPROCESSABLE_ENTITY, "UPLOADSESSION_EXPIRATION_TOO_FAR"),
Self::SessionExpired(_) => (StatusCode::UNPROCESSABLE_ENTITY, "UPLOADSESSION_SESSION_EXPIRED"),
};
let body = serde_json::json!({
"success": false,
"error": code,
"message": self.to_string(),
});
(status, Json(body)).into_response()
}
}
pub mod upload_session_errors {
pub const BUCKET_NOT_FOUND: &str = "UPLOADSESSION_BUCKET_NOT_FOUND";
pub const BUCKET_NOT_ACTIVE: &str = "UPLOADSESSION_BUCKET_NOT_ACTIVE";
pub const INVALID_FILE_SIZE: &str = "UPLOADSESSION_INVALID_FILE_SIZE";
pub const CHUNK_TOO_SMALL: &str = "UPLOADSESSION_CHUNK_TOO_SMALL";
pub const CHUNK_TOO_LARGE: &str = "UPLOADSESSION_CHUNK_TOO_LARGE";
pub const INVALID_CHUNK_COUNT: &str = "UPLOADSESSION_INVALID_CHUNK_COUNT";
pub const INVALID_PATH: &str = "UPLOADSESSION_INVALID_PATH";
pub const INVALID_EXPIRATION: &str = "UPLOADSESSION_INVALID_EXPIRATION";
pub const EXPIRATION_TOO_FAR: &str = "UPLOADSESSION_EXPIRATION_TOO_FAR";
pub const SESSION_EXPIRED: &str = "UPLOADSESSION_SESSION_EXPIRED";
}
pub fn create_upload_session_routes(service: Arc<UploadSessionService>) -> Router {
BackboneCrudHandler::<UploadSessionService, UploadSession, CreateUploadSessionDto, UpdateUploadSessionDto, UploadSessionResponseDto>::routes(
service,
"/upload_sessions",
)
}
pub fn create_upload_session_read_routes(service: Arc<UploadSessionService>) -> Router {
BackboneCrudHandler::<UploadSessionService, UploadSession, CreateUploadSessionDto, UpdateUploadSessionDto, UploadSessionResponseDto>::read_routes(
service,
"/upload_sessions",
)
}
pub fn create_upload_session_write_routes(service: Arc<UploadSessionService>) -> Router {
BackboneCrudHandler::<UploadSessionService, UploadSession, CreateUploadSessionDto, UpdateUploadSessionDto, UploadSessionResponseDto>::write_routes(
service,
"/upload_sessions",
)
}
#[cfg(feature = "auth")]
pub fn create_protected_upload_session_routes<A: AuthMiddleware + Send + Sync + 'static>(
service: Arc<UploadSessionService>,
auth: Arc<A>,
) -> Router {
use axum::middleware;
use axum::response::IntoResponse;
let auth_layer = auth.clone();
create_upload_session_routes(service)
.layer(middleware::from_fn(move |mut req: axum::extract::Request, next: axum::middleware::Next| {
let auth = auth_layer.clone();
async move {
let token = req.headers()
.get(axum::http::header::AUTHORIZATION)
.and_then(|h| h.to_str().ok())
.and_then(|raw| raw.strip_prefix("Bearer ").or_else(|| raw.strip_prefix("bearer ")))
.unwrap_or("");
match auth.authenticate(token).await {
Ok(ctx) => {
req.extensions_mut().insert(ctx);
next.run(req).await
}
Err(_) => {
(axum::http::StatusCode::UNAUTHORIZED,
axum::Json(serde_json::json!({
"success": false,
"error": "unauthorized",
"message": "Authentication required"
}))
).into_response()
}
}
}
}))
}
pub async fn start_upload_transition(
axum::extract::State(service): axum::extract::State<Arc<UploadSessionService>>,
axum::extract::Path(id): axum::extract::Path<String>,
#[cfg(feature = "auth")] axum::Extension(auth): axum::Extension<AuthContext>,
) -> impl axum::response::IntoResponse {
use axum::{http::StatusCode, Json};
let entity = match service.get_by_id(&id).await {
Ok(Some(e)) => e,
Ok(None) => return (StatusCode::NOT_FOUND, Json(ApiResponse::<UploadSessionResponseDto>::not_found("UploadSession", &id))),
Err(e) => return (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<UploadSessionResponseDto>::error(e.to_string()))),
};
#[cfg(feature = "auth")]
{
let allowed_roles = UploadSessionTransition::StartUpload.allowed_roles();
let has_specific_perm = auth.permissions.iter().any(|p| p == "upload_session:transition:start_upload");
let has_update_perm = auth.permissions.iter().any(|p| p == "upload_session:update");
let has_role = auth.roles.iter().any(|r| allowed_roles.contains(&r.as_str()));
if !has_specific_perm && !has_update_perm && !has_role {
return (StatusCode::FORBIDDEN, Json(ApiResponse::<UploadSessionResponseDto>::error("Insufficient permissions for start_upload transition")));
}
}
let current_state: UploadSessionState = entity.status.to_string().parse()
.unwrap_or(UploadSessionState::default());
let sm = UploadSessionStateMachine::from_state(current_state);
if !sm.can_transition(UploadSessionTransition::StartUpload) {
return (StatusCode::BAD_REQUEST, Json(ApiResponse::<UploadSessionResponseDto>::error("Transition not allowed from current state")));
}
let mut fields: HashMap<String, serde_json::Value> = HashMap::new();
fields.insert("status".to_string(), serde_json::Value::String("Uploading".to_string()));
match service.partial_update(&id, fields).await {
Ok(Some(updated)) => {
let response: UploadSessionResponseDto = updated.into();
(StatusCode::OK, Json(ApiResponse::ok(response)))
}
Ok(None) => (StatusCode::NOT_FOUND, Json(ApiResponse::<UploadSessionResponseDto>::not_found("UploadSession", &id))),
Err(e) => (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<UploadSessionResponseDto>::error(e.to_string()))),
}
}
pub async fn add_part_transition(
axum::extract::State(service): axum::extract::State<Arc<UploadSessionService>>,
axum::extract::Path(id): axum::extract::Path<String>,
#[cfg(feature = "auth")] axum::Extension(auth): axum::Extension<AuthContext>,
) -> impl axum::response::IntoResponse {
use axum::{http::StatusCode, Json};
let entity = match service.get_by_id(&id).await {
Ok(Some(e)) => e,
Ok(None) => return (StatusCode::NOT_FOUND, Json(ApiResponse::<UploadSessionResponseDto>::not_found("UploadSession", &id))),
Err(e) => return (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<UploadSessionResponseDto>::error(e.to_string()))),
};
#[cfg(feature = "auth")]
{
let allowed_roles = UploadSessionTransition::AddPart.allowed_roles();
let has_specific_perm = auth.permissions.iter().any(|p| p == "upload_session:transition:add_part");
let has_update_perm = auth.permissions.iter().any(|p| p == "upload_session:update");
let has_role = auth.roles.iter().any(|r| allowed_roles.contains(&r.as_str()));
if !has_specific_perm && !has_update_perm && !has_role {
return (StatusCode::FORBIDDEN, Json(ApiResponse::<UploadSessionResponseDto>::error("Insufficient permissions for add_part transition")));
}
}
let current_state: UploadSessionState = entity.status.to_string().parse()
.unwrap_or(UploadSessionState::default());
let sm = UploadSessionStateMachine::from_state(current_state);
if !sm.can_transition(UploadSessionTransition::AddPart) {
return (StatusCode::BAD_REQUEST, Json(ApiResponse::<UploadSessionResponseDto>::error("Transition not allowed from current state")));
}
let mut fields: HashMap<String, serde_json::Value> = HashMap::new();
fields.insert("status".to_string(), serde_json::Value::String("Uploading".to_string()));
match service.partial_update(&id, fields).await {
Ok(Some(updated)) => {
let response: UploadSessionResponseDto = updated.into();
(StatusCode::OK, Json(ApiResponse::ok(response)))
}
Ok(None) => (StatusCode::NOT_FOUND, Json(ApiResponse::<UploadSessionResponseDto>::not_found("UploadSession", &id))),
Err(e) => (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<UploadSessionResponseDto>::error(e.to_string()))),
}
}
pub async fn complete_transition(
axum::extract::State(service): axum::extract::State<Arc<UploadSessionService>>,
axum::extract::Path(id): axum::extract::Path<String>,
#[cfg(feature = "auth")] axum::Extension(auth): axum::Extension<AuthContext>,
) -> impl axum::response::IntoResponse {
use axum::{http::StatusCode, Json};
let entity = match service.get_by_id(&id).await {
Ok(Some(e)) => e,
Ok(None) => return (StatusCode::NOT_FOUND, Json(ApiResponse::<UploadSessionResponseDto>::not_found("UploadSession", &id))),
Err(e) => return (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<UploadSessionResponseDto>::error(e.to_string()))),
};
#[cfg(feature = "auth")]
{
let allowed_roles = UploadSessionTransition::Complete.allowed_roles();
let has_specific_perm = auth.permissions.iter().any(|p| p == "upload_session:transition:complete");
let has_update_perm = auth.permissions.iter().any(|p| p == "upload_session:update");
let has_role = auth.roles.iter().any(|r| allowed_roles.contains(&r.as_str()));
if !has_specific_perm && !has_update_perm && !has_role {
return (StatusCode::FORBIDDEN, Json(ApiResponse::<UploadSessionResponseDto>::error("Insufficient permissions for complete transition")));
}
}
let current_state: UploadSessionState = entity.status.to_string().parse()
.unwrap_or(UploadSessionState::default());
let sm = UploadSessionStateMachine::from_state(current_state);
if !sm.can_transition(UploadSessionTransition::Complete) {
return (StatusCode::BAD_REQUEST, Json(ApiResponse::<UploadSessionResponseDto>::error("Transition not allowed from current state")));
}
let mut fields: HashMap<String, serde_json::Value> = HashMap::new();
fields.insert("status".to_string(), serde_json::Value::String("Completing".to_string()));
match service.partial_update(&id, fields).await {
Ok(Some(updated)) => {
let response: UploadSessionResponseDto = updated.into();
(StatusCode::OK, Json(ApiResponse::ok(response)))
}
Ok(None) => (StatusCode::NOT_FOUND, Json(ApiResponse::<UploadSessionResponseDto>::not_found("UploadSession", &id))),
Err(e) => (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<UploadSessionResponseDto>::error(e.to_string()))),
}
}
pub async fn finish_transition(
axum::extract::State(service): axum::extract::State<Arc<UploadSessionService>>,
axum::extract::Path(id): axum::extract::Path<String>,
#[cfg(feature = "auth")] axum::Extension(auth): axum::Extension<AuthContext>,
) -> impl axum::response::IntoResponse {
use axum::{http::StatusCode, Json};
let entity = match service.get_by_id(&id).await {
Ok(Some(e)) => e,
Ok(None) => return (StatusCode::NOT_FOUND, Json(ApiResponse::<UploadSessionResponseDto>::not_found("UploadSession", &id))),
Err(e) => return (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<UploadSessionResponseDto>::error(e.to_string()))),
};
#[cfg(feature = "auth")]
{
let allowed_roles = UploadSessionTransition::Finish.allowed_roles();
let has_specific_perm = auth.permissions.iter().any(|p| p == "upload_session:transition:finish");
let has_update_perm = auth.permissions.iter().any(|p| p == "upload_session:update");
let has_role = auth.roles.iter().any(|r| allowed_roles.contains(&r.as_str()));
if !has_specific_perm && !has_update_perm && !has_role {
return (StatusCode::FORBIDDEN, Json(ApiResponse::<UploadSessionResponseDto>::error("Insufficient permissions for finish transition")));
}
}
let current_state: UploadSessionState = entity.status.to_string().parse()
.unwrap_or(UploadSessionState::default());
let sm = UploadSessionStateMachine::from_state(current_state);
if !sm.can_transition(UploadSessionTransition::Finish) {
return (StatusCode::BAD_REQUEST, Json(ApiResponse::<UploadSessionResponseDto>::error("Transition not allowed from current state")));
}
let mut fields: HashMap<String, serde_json::Value> = HashMap::new();
fields.insert("status".to_string(), serde_json::Value::String("Completed".to_string()));
match service.partial_update(&id, fields).await {
Ok(Some(updated)) => {
let response: UploadSessionResponseDto = updated.into();
(StatusCode::OK, Json(ApiResponse::ok(response)))
}
Ok(None) => (StatusCode::NOT_FOUND, Json(ApiResponse::<UploadSessionResponseDto>::not_found("UploadSession", &id))),
Err(e) => (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<UploadSessionResponseDto>::error(e.to_string()))),
}
}
pub async fn fail_transition(
axum::extract::State(service): axum::extract::State<Arc<UploadSessionService>>,
axum::extract::Path(id): axum::extract::Path<String>,
#[cfg(feature = "auth")] axum::Extension(auth): axum::Extension<AuthContext>,
) -> impl axum::response::IntoResponse {
use axum::{http::StatusCode, Json};
let entity = match service.get_by_id(&id).await {
Ok(Some(e)) => e,
Ok(None) => return (StatusCode::NOT_FOUND, Json(ApiResponse::<UploadSessionResponseDto>::not_found("UploadSession", &id))),
Err(e) => return (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<UploadSessionResponseDto>::error(e.to_string()))),
};
#[cfg(feature = "auth")]
{
let allowed_roles = UploadSessionTransition::Fail.allowed_roles();
let has_specific_perm = auth.permissions.iter().any(|p| p == "upload_session:transition:fail");
let has_update_perm = auth.permissions.iter().any(|p| p == "upload_session:update");
let has_role = auth.roles.iter().any(|r| allowed_roles.contains(&r.as_str()));
if !has_specific_perm && !has_update_perm && !has_role {
return (StatusCode::FORBIDDEN, Json(ApiResponse::<UploadSessionResponseDto>::error("Insufficient permissions for fail transition")));
}
}
let current_state: UploadSessionState = entity.status.to_string().parse()
.unwrap_or(UploadSessionState::default());
let sm = UploadSessionStateMachine::from_state(current_state);
if !sm.can_transition(UploadSessionTransition::Fail) {
return (StatusCode::BAD_REQUEST, Json(ApiResponse::<UploadSessionResponseDto>::error("Transition not allowed from current state")));
}
let mut fields: HashMap<String, serde_json::Value> = HashMap::new();
fields.insert("status".to_string(), serde_json::Value::String("Failed".to_string()));
match service.partial_update(&id, fields).await {
Ok(Some(updated)) => {
let response: UploadSessionResponseDto = updated.into();
(StatusCode::OK, Json(ApiResponse::ok(response)))
}
Ok(None) => (StatusCode::NOT_FOUND, Json(ApiResponse::<UploadSessionResponseDto>::not_found("UploadSession", &id))),
Err(e) => (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<UploadSessionResponseDto>::error(e.to_string()))),
}
}
pub async fn abort_transition(
axum::extract::State(service): axum::extract::State<Arc<UploadSessionService>>,
axum::extract::Path(id): axum::extract::Path<String>,
#[cfg(feature = "auth")] axum::Extension(auth): axum::Extension<AuthContext>,
) -> impl axum::response::IntoResponse {
use axum::{http::StatusCode, Json};
let entity = match service.get_by_id(&id).await {
Ok(Some(e)) => e,
Ok(None) => return (StatusCode::NOT_FOUND, Json(ApiResponse::<UploadSessionResponseDto>::not_found("UploadSession", &id))),
Err(e) => return (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<UploadSessionResponseDto>::error(e.to_string()))),
};
#[cfg(feature = "auth")]
{
let allowed_roles = UploadSessionTransition::Abort.allowed_roles();
let has_specific_perm = auth.permissions.iter().any(|p| p == "upload_session:transition:abort");
let has_update_perm = auth.permissions.iter().any(|p| p == "upload_session:update");
let has_role = auth.roles.iter().any(|r| allowed_roles.contains(&r.as_str()));
if !has_specific_perm && !has_update_perm && !has_role {
return (StatusCode::FORBIDDEN, Json(ApiResponse::<UploadSessionResponseDto>::error("Insufficient permissions for abort transition")));
}
}
let current_state: UploadSessionState = entity.status.to_string().parse()
.unwrap_or(UploadSessionState::default());
let sm = UploadSessionStateMachine::from_state(current_state);
if !sm.can_transition(UploadSessionTransition::Abort) {
return (StatusCode::BAD_REQUEST, Json(ApiResponse::<UploadSessionResponseDto>::error("Transition not allowed from current state")));
}
let mut fields: HashMap<String, serde_json::Value> = HashMap::new();
fields.insert("status".to_string(), serde_json::Value::String("Failed".to_string()));
match service.partial_update(&id, fields).await {
Ok(Some(updated)) => {
let response: UploadSessionResponseDto = updated.into();
(StatusCode::OK, Json(ApiResponse::ok(response)))
}
Ok(None) => (StatusCode::NOT_FOUND, Json(ApiResponse::<UploadSessionResponseDto>::not_found("UploadSession", &id))),
Err(e) => (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<UploadSessionResponseDto>::error(e.to_string()))),
}
}
pub async fn expire_transition(
axum::extract::State(service): axum::extract::State<Arc<UploadSessionService>>,
axum::extract::Path(id): axum::extract::Path<String>,
#[cfg(feature = "auth")] axum::Extension(auth): axum::Extension<AuthContext>,
) -> impl axum::response::IntoResponse {
use axum::{http::StatusCode, Json};
let entity = match service.get_by_id(&id).await {
Ok(Some(e)) => e,
Ok(None) => return (StatusCode::NOT_FOUND, Json(ApiResponse::<UploadSessionResponseDto>::not_found("UploadSession", &id))),
Err(e) => return (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<UploadSessionResponseDto>::error(e.to_string()))),
};
#[cfg(feature = "auth")]
{
let allowed_roles = UploadSessionTransition::Expire.allowed_roles();
let has_specific_perm = auth.permissions.iter().any(|p| p == "upload_session:transition:expire");
let has_update_perm = auth.permissions.iter().any(|p| p == "upload_session:update");
let has_role = auth.roles.iter().any(|r| allowed_roles.contains(&r.as_str()));
if !has_specific_perm && !has_update_perm && !has_role {
return (StatusCode::FORBIDDEN, Json(ApiResponse::<UploadSessionResponseDto>::error("Insufficient permissions for expire transition")));
}
}
let current_state: UploadSessionState = entity.status.to_string().parse()
.unwrap_or(UploadSessionState::default());
let sm = UploadSessionStateMachine::from_state(current_state);
if !sm.can_transition(UploadSessionTransition::Expire) {
return (StatusCode::BAD_REQUEST, Json(ApiResponse::<UploadSessionResponseDto>::error("Transition not allowed from current state")));
}
let mut fields: HashMap<String, serde_json::Value> = HashMap::new();
fields.insert("status".to_string(), serde_json::Value::String("Expired".to_string()));
match service.partial_update(&id, fields).await {
Ok(Some(updated)) => {
let response: UploadSessionResponseDto = updated.into();
(StatusCode::OK, Json(ApiResponse::ok(response)))
}
Ok(None) => (StatusCode::NOT_FOUND, Json(ApiResponse::<UploadSessionResponseDto>::not_found("UploadSession", &id))),
Err(e) => (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<UploadSessionResponseDto>::error(e.to_string()))),
}
}
pub async fn retry_transition(
axum::extract::State(service): axum::extract::State<Arc<UploadSessionService>>,
axum::extract::Path(id): axum::extract::Path<String>,
#[cfg(feature = "auth")] axum::Extension(auth): axum::Extension<AuthContext>,
) -> impl axum::response::IntoResponse {
use axum::{http::StatusCode, Json};
let entity = match service.get_by_id(&id).await {
Ok(Some(e)) => e,
Ok(None) => return (StatusCode::NOT_FOUND, Json(ApiResponse::<UploadSessionResponseDto>::not_found("UploadSession", &id))),
Err(e) => return (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<UploadSessionResponseDto>::error(e.to_string()))),
};
#[cfg(feature = "auth")]
{
let allowed_roles = UploadSessionTransition::Retry.allowed_roles();
let has_specific_perm = auth.permissions.iter().any(|p| p == "upload_session:transition:retry");
let has_update_perm = auth.permissions.iter().any(|p| p == "upload_session:update");
let has_role = auth.roles.iter().any(|r| allowed_roles.contains(&r.as_str()));
if !has_specific_perm && !has_update_perm && !has_role {
return (StatusCode::FORBIDDEN, Json(ApiResponse::<UploadSessionResponseDto>::error("Insufficient permissions for retry transition")));
}
}
let current_state: UploadSessionState = entity.status.to_string().parse()
.unwrap_or(UploadSessionState::default());
let sm = UploadSessionStateMachine::from_state(current_state);
if !sm.can_transition(UploadSessionTransition::Retry) {
return (StatusCode::BAD_REQUEST, Json(ApiResponse::<UploadSessionResponseDto>::error("Transition not allowed from current state")));
}
let mut fields: HashMap<String, serde_json::Value> = HashMap::new();
fields.insert("status".to_string(), serde_json::Value::String("Initiated".to_string()));
match service.partial_update(&id, fields).await {
Ok(Some(updated)) => {
let response: UploadSessionResponseDto = updated.into();
(StatusCode::OK, Json(ApiResponse::ok(response)))
}
Ok(None) => (StatusCode::NOT_FOUND, Json(ApiResponse::<UploadSessionResponseDto>::not_found("UploadSession", &id))),
Err(e) => (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<UploadSessionResponseDto>::error(e.to_string()))),
}
}
pub fn create_upload_session_transition_routes(service: Arc<UploadSessionService>) -> Router {
use axum::routing::post;
Router::new()
.route("/upload_sessions/:id/transitions/start_upload", post(start_upload_transition))
.route("/upload_sessions/:id/transitions/add_part", post(add_part_transition))
.route("/upload_sessions/:id/transitions/complete", post(complete_transition))
.route("/upload_sessions/:id/transitions/finish", post(finish_transition))
.route("/upload_sessions/:id/transitions/fail", post(fail_transition))
.route("/upload_sessions/:id/transitions/abort", post(abort_transition))
.route("/upload_sessions/:id/transitions/expire", post(expire_transition))
.route("/upload_sessions/:id/transitions/retry", post(retry_transition))
.with_state(service)
}