Skip to main content

backbone_bucket/presentation/http/
processing_job_handler.rs

1//! ProcessingJob REST handlers
2//!
3//! Generated by metaphor-schema. Do not edit manually.
4//!
5//! Uses Axum and backbone-core's BackboneCrudHandler for all 12 CRUD endpoints.
6
7use std::collections::HashMap;
8use std::sync::Arc;
9
10use axum::Router;
11use serde::{Deserialize, Serialize};
12use uuid::Uuid;
13use chrono::{DateTime, Utc};
14
15// Backbone framework imports
16use backbone_core::http::{ApiResponse, BackboneCrudHandler};
17
18// Auth integration (optional)
19#[cfg(feature = "auth")]
20use backbone_auth::middleware::AuthContext;
21#[cfg(feature = "auth")]
22use backbone_auth::AuthMiddleware;
23
24// Domain imports
25use crate::domain::entity::*;
26use crate::application::service::{ProcessingJobService, ServiceError};
27
28// DTO imports
29use crate::presentation::dto::{CreateProcessingJobDto, UpdateProcessingJobDto, PatchProcessingJobDto, ProcessingJobResponseDto};
30
31use crate::domain::state_machine::{ProcessingJobState, ProcessingJobStateMachine, ProcessingJobTransition};
32
33/// Application error type
34#[derive(Debug, thiserror::Error)]
35pub enum ProcessingJobError {
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    // Domain-specific errors from hook rules
45    #[error("File must exist: {0}")]
46    FileNotFound(String),
47    #[error("File must be accessible: {0}")]
48    FileNotAccessible(String),
49    #[error("Job type must be supported: {0}")]
50    InvalidJobType(String),
51    #[error("Job type must match file type: {0}")]
52    JobTypeMismatch(String),
53    #[error("A similar job is already pending: {0}")]
54    DuplicateJob(String),
55    #[error("Priority must be 0-10: {0}")]
56    InvalidPriority(String),
57    #[error("Cannot modify completed jobs: {0}")]
58    JobCompleted(String),
59}
60
61impl From<ServiceError> for ProcessingJobError {
62    fn from(err: ServiceError) -> Self {
63        match err {
64            ServiceError::NotFound => Self::NotFound(err.to_string()),
65            ServiceError::Validation(ref msg) => Self::Validation(msg.clone()),
66            ServiceError::AlreadyExists(ref msg) => Self::Validation(msg.clone()),
67            ServiceError::Repository(ref e) => Self::Database(e.to_string()),
68            ServiceError::Internal(ref msg) => Self::Internal(msg.clone()),
69            ServiceError::Violations(_) => Self::Validation(err.to_string()),
70        }
71    }
72}
73
74impl axum::response::IntoResponse for ProcessingJobError {
75    fn into_response(self) -> axum::response::Response {
76        use axum::http::StatusCode;
77        use axum::Json;
78
79        let (status, code) = match &self {
80            Self::NotFound(_) => (StatusCode::NOT_FOUND, "PROCESSINGJOB_NOT_FOUND"),
81            Self::Validation(_) => (StatusCode::BAD_REQUEST, "PROCESSINGJOB_VALIDATION_ERROR"),
82            Self::Database(_) => (StatusCode::INTERNAL_SERVER_ERROR, "PROCESSINGJOB_DATABASE_ERROR"),
83            Self::Internal(_) => (StatusCode::INTERNAL_SERVER_ERROR, "PROCESSINGJOB_INTERNAL_ERROR"),
84            Self::FileNotFound(_) => (StatusCode::UNPROCESSABLE_ENTITY, "PROCESSINGJOB_FILE_NOT_FOUND"),
85            Self::FileNotAccessible(_) => (StatusCode::UNPROCESSABLE_ENTITY, "PROCESSINGJOB_FILE_NOT_ACCESSIBLE"),
86            Self::InvalidJobType(_) => (StatusCode::UNPROCESSABLE_ENTITY, "PROCESSINGJOB_INVALID_JOB_TYPE"),
87            Self::JobTypeMismatch(_) => (StatusCode::UNPROCESSABLE_ENTITY, "PROCESSINGJOB_JOB_TYPE_MISMATCH"),
88            Self::DuplicateJob(_) => (StatusCode::UNPROCESSABLE_ENTITY, "PROCESSINGJOB_DUPLICATE_JOB"),
89            Self::InvalidPriority(_) => (StatusCode::UNPROCESSABLE_ENTITY, "PROCESSINGJOB_INVALID_PRIORITY"),
90            Self::JobCompleted(_) => (StatusCode::UNPROCESSABLE_ENTITY, "PROCESSINGJOB_JOB_COMPLETED"),
91        };
92
93        let body = serde_json::json!({
94            "success": false,
95            "error": code,
96            "message": self.to_string(),
97        });
98
99        (status, Json(body)).into_response()
100    }
101}
102
103/// Domain-specific error codes for ProcessingJob
104pub mod processing_job_errors {
105    pub const FILE_NOT_FOUND: &str = "PROCESSINGJOB_FILE_NOT_FOUND";
106    pub const FILE_NOT_ACCESSIBLE: &str = "PROCESSINGJOB_FILE_NOT_ACCESSIBLE";
107    pub const INVALID_JOB_TYPE: &str = "PROCESSINGJOB_INVALID_JOB_TYPE";
108    pub const JOB_TYPE_MISMATCH: &str = "PROCESSINGJOB_JOB_TYPE_MISMATCH";
109    pub const DUPLICATE_JOB: &str = "PROCESSINGJOB_DUPLICATE_JOB";
110    pub const INVALID_PRIORITY: &str = "PROCESSINGJOB_INVALID_PRIORITY";
111    pub const JOB_COMPLETED: &str = "PROCESSINGJOB_JOB_COMPLETED";
112}
113
114// =============================================================================
115// Route Configuration
116// =============================================================================
117
118/// Create Axum router with all 16 Backbone endpoints for ProcessingJob.
119///
120/// # Routes
121///
122/// | Method | Path | Description |
123/// |--------|------|-------------|
124/// | GET | /processing_jobs | List with pagination |
125/// | POST | /processing_jobs | Create new |
126/// | GET | /processing_jobs/:id | Get by ID |
127/// | PUT | /processing_jobs/:id | Full update |
128/// | PATCH | /processing_jobs/:id | Partial update |
129/// | DELETE | /processing_jobs/:id | Soft delete |
130/// | POST | /processing_jobs/bulk | Bulk create |
131/// | POST | /processing_jobs/upsert | Upsert |
132/// | GET | /processing_jobs/trash | List deleted |
133/// | POST | /processing_jobs/:id/restore | Restore |
134/// | DELETE | /processing_jobs/empty | Empty trash |
135/// | GET | /processing_jobs/:id/deleted | Get deleted by ID |
136/// | DELETE | /processing_jobs/trash/:id | Permanent delete from trash |
137/// | GET | /processing_jobs/count | Count active entities |
138/// | GET | /processing_jobs/trash/count | Count deleted entities |
139///
140/// # Example
141///
142/// ```text
143/// let service = Arc::new(ProcessingJobService::with_repository(repository));
144/// let router = create_processing_job_routes(service);
145/// ```
146pub fn create_processing_job_routes(service: Arc<ProcessingJobService>) -> Router {
147    BackboneCrudHandler::<ProcessingJobService, ProcessingJob, CreateProcessingJobDto, UpdateProcessingJobDto, ProcessingJobResponseDto>::routes(
148        service,
149        "/processing_jobs",
150    )
151}
152
153/// Create Axum router with only the read (GET) endpoints for ProcessingJob.
154///
155/// Safe for public, unauthenticated exposure (e.g., reference data).
156/// Mutations must be served separately via `create_processing_job_write_routes`,
157/// typically wrapped in an auth middleware layer.
158pub fn create_processing_job_read_routes(service: Arc<ProcessingJobService>) -> Router {
159    BackboneCrudHandler::<ProcessingJobService, ProcessingJob, CreateProcessingJobDto, UpdateProcessingJobDto, ProcessingJobResponseDto>::read_routes(
160        service,
161        "/processing_jobs",
162    )
163}
164
165/// Create Axum router with only the write (mutation) endpoints for ProcessingJob.
166///
167/// These routes must NOT be publicly exposed. Wrap them with an auth
168/// middleware before nesting into the application router.
169///
170/// # This is unguarded generic CRUD, not a validated write path
171///
172/// These are plain create/update/patch/delete mutations over the entity row —
173/// they bypass all business invariants. If the module exposes a validated write
174/// service (e.g. a command router over its domain engine), serve THAT instead
175/// for any mutation that must respect domain rules.
176pub fn create_processing_job_write_routes(service: Arc<ProcessingJobService>) -> Router {
177    BackboneCrudHandler::<ProcessingJobService, ProcessingJob, CreateProcessingJobDto, UpdateProcessingJobDto, ProcessingJobResponseDto>::write_routes(
178        service,
179        "/processing_jobs",
180    )
181}
182
183/// Create authenticated routes with auth middleware.
184///
185/// Requires the `auth` feature flag. The `AuthMiddleware` implementation
186/// is responsible for extracting and validating tokens, then providing
187/// an `AuthContext` via request extensions.
188#[cfg(feature = "auth")]
189pub fn create_protected_processing_job_routes<A: AuthMiddleware + Send + Sync + 'static>(
190    service: Arc<ProcessingJobService>,
191    auth: Arc<A>,
192) -> Router {
193    use axum::middleware;
194    use axum::response::IntoResponse;
195
196    let auth_layer = auth.clone();
197    create_processing_job_routes(service)
198        .layer(middleware::from_fn(move |mut req: axum::extract::Request, next: axum::middleware::Next| {
199            let auth = auth_layer.clone();
200            async move {
201                let token = req.headers()
202                    .get(axum::http::header::AUTHORIZATION)
203                    .and_then(|h| h.to_str().ok())
204                    .and_then(|raw| raw.strip_prefix("Bearer ").or_else(|| raw.strip_prefix("bearer ")))
205                    .unwrap_or("");
206                match auth.authenticate(token).await {
207                    Ok(ctx) => {
208                        req.extensions_mut().insert(ctx);
209                        next.run(req).await
210                    }
211                    Err(_) => {
212                        (axum::http::StatusCode::UNAUTHORIZED,
213                         axum::Json(serde_json::json!({
214                             "success": false,
215                             "error": "unauthorized",
216                             "message": "Authentication required"
217                         }))
218                        ).into_response()
219                    }
220                }
221            }
222        }))
223}
224
225// =============================================================================
226// State Transition Handlers
227// =============================================================================
228
229/// Execute start transition on a ProcessingJob.
230///
231/// POST /processing_jobs/:id/transitions/start
232pub async fn start_transition(
233    axum::extract::State(service): axum::extract::State<Arc<ProcessingJobService>>,
234    axum::extract::Path(id): axum::extract::Path<String>,
235    #[cfg(feature = "auth")] axum::Extension(auth): axum::Extension<AuthContext>,
236) -> impl axum::response::IntoResponse {
237    use axum::{http::StatusCode, Json};
238
239    // Get current entity
240    let entity = match service.get_by_id(&id).await {
241        Ok(Some(e)) => e,
242        Ok(None) => return (StatusCode::NOT_FOUND, Json(ApiResponse::<ProcessingJobResponseDto>::not_found("ProcessingJob", &id))),
243        Err(e) => return (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<ProcessingJobResponseDto>::error(e.to_string()))),
244    };
245
246    // Check permission (if auth enabled)
247    #[cfg(feature = "auth")]
248    {
249        let allowed_roles = ProcessingJobTransition::Start.allowed_roles();
250        let has_specific_perm = auth.permissions.iter().any(|p| p == "processing_job:transition:start");
251        let has_update_perm = auth.permissions.iter().any(|p| p == "processing_job:update");
252        let has_role = auth.roles.iter().any(|r| allowed_roles.contains(&r.as_str()));
253        if !has_specific_perm && !has_update_perm && !has_role {
254            return (StatusCode::FORBIDDEN, Json(ApiResponse::<ProcessingJobResponseDto>::error("Insufficient permissions for start transition")));
255        }
256    }
257
258    // Create state machine from entity's actual status and validate transition
259    let current_state: ProcessingJobState = entity.status.to_string().parse()
260        .unwrap_or(ProcessingJobState::default());
261    let sm = ProcessingJobStateMachine::from_state(current_state);
262    if !sm.can_transition(ProcessingJobTransition::Start) {
263        return (StatusCode::BAD_REQUEST, Json(ApiResponse::<ProcessingJobResponseDto>::error("Transition not allowed from current state")));
264    }
265
266    // Apply transition via partial update
267    let mut fields: HashMap<String, serde_json::Value> = HashMap::new();
268    fields.insert("status".to_string(), serde_json::Value::String("Running".to_string()));
269
270    match service.partial_update(&id, fields).await {
271        Ok(Some(updated)) => {
272            let response: ProcessingJobResponseDto = updated.into();
273            (StatusCode::OK, Json(ApiResponse::ok(response)))
274        }
275        Ok(None) => (StatusCode::NOT_FOUND, Json(ApiResponse::<ProcessingJobResponseDto>::not_found("ProcessingJob", &id))),
276        Err(e) => (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<ProcessingJobResponseDto>::error(e.to_string()))),
277    }
278}
279
280/// Execute complete transition on a ProcessingJob.
281///
282/// POST /processing_jobs/:id/transitions/complete
283pub async fn complete_transition(
284    axum::extract::State(service): axum::extract::State<Arc<ProcessingJobService>>,
285    axum::extract::Path(id): axum::extract::Path<String>,
286    #[cfg(feature = "auth")] axum::Extension(auth): axum::Extension<AuthContext>,
287) -> impl axum::response::IntoResponse {
288    use axum::{http::StatusCode, Json};
289
290    // Get current entity
291    let entity = match service.get_by_id(&id).await {
292        Ok(Some(e)) => e,
293        Ok(None) => return (StatusCode::NOT_FOUND, Json(ApiResponse::<ProcessingJobResponseDto>::not_found("ProcessingJob", &id))),
294        Err(e) => return (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<ProcessingJobResponseDto>::error(e.to_string()))),
295    };
296
297    // Check permission (if auth enabled)
298    #[cfg(feature = "auth")]
299    {
300        let allowed_roles = ProcessingJobTransition::Complete.allowed_roles();
301        let has_specific_perm = auth.permissions.iter().any(|p| p == "processing_job:transition:complete");
302        let has_update_perm = auth.permissions.iter().any(|p| p == "processing_job:update");
303        let has_role = auth.roles.iter().any(|r| allowed_roles.contains(&r.as_str()));
304        if !has_specific_perm && !has_update_perm && !has_role {
305            return (StatusCode::FORBIDDEN, Json(ApiResponse::<ProcessingJobResponseDto>::error("Insufficient permissions for complete transition")));
306        }
307    }
308
309    // Create state machine from entity's actual status and validate transition
310    let current_state: ProcessingJobState = entity.status.to_string().parse()
311        .unwrap_or(ProcessingJobState::default());
312    let sm = ProcessingJobStateMachine::from_state(current_state);
313    if !sm.can_transition(ProcessingJobTransition::Complete) {
314        return (StatusCode::BAD_REQUEST, Json(ApiResponse::<ProcessingJobResponseDto>::error("Transition not allowed from current state")));
315    }
316
317    // Apply transition via partial update
318    let mut fields: HashMap<String, serde_json::Value> = HashMap::new();
319    fields.insert("status".to_string(), serde_json::Value::String("Completed".to_string()));
320
321    match service.partial_update(&id, fields).await {
322        Ok(Some(updated)) => {
323            let response: ProcessingJobResponseDto = updated.into();
324            (StatusCode::OK, Json(ApiResponse::ok(response)))
325        }
326        Ok(None) => (StatusCode::NOT_FOUND, Json(ApiResponse::<ProcessingJobResponseDto>::not_found("ProcessingJob", &id))),
327        Err(e) => (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<ProcessingJobResponseDto>::error(e.to_string()))),
328    }
329}
330
331/// Execute fail transition on a ProcessingJob.
332///
333/// POST /processing_jobs/:id/transitions/fail
334pub async fn fail_transition(
335    axum::extract::State(service): axum::extract::State<Arc<ProcessingJobService>>,
336    axum::extract::Path(id): axum::extract::Path<String>,
337    #[cfg(feature = "auth")] axum::Extension(auth): axum::Extension<AuthContext>,
338) -> impl axum::response::IntoResponse {
339    use axum::{http::StatusCode, Json};
340
341    // Get current entity
342    let entity = match service.get_by_id(&id).await {
343        Ok(Some(e)) => e,
344        Ok(None) => return (StatusCode::NOT_FOUND, Json(ApiResponse::<ProcessingJobResponseDto>::not_found("ProcessingJob", &id))),
345        Err(e) => return (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<ProcessingJobResponseDto>::error(e.to_string()))),
346    };
347
348    // Check permission (if auth enabled)
349    #[cfg(feature = "auth")]
350    {
351        let allowed_roles = ProcessingJobTransition::Fail.allowed_roles();
352        let has_specific_perm = auth.permissions.iter().any(|p| p == "processing_job:transition:fail");
353        let has_update_perm = auth.permissions.iter().any(|p| p == "processing_job:update");
354        let has_role = auth.roles.iter().any(|r| allowed_roles.contains(&r.as_str()));
355        if !has_specific_perm && !has_update_perm && !has_role {
356            return (StatusCode::FORBIDDEN, Json(ApiResponse::<ProcessingJobResponseDto>::error("Insufficient permissions for fail transition")));
357        }
358    }
359
360    // Create state machine from entity's actual status and validate transition
361    let current_state: ProcessingJobState = entity.status.to_string().parse()
362        .unwrap_or(ProcessingJobState::default());
363    let sm = ProcessingJobStateMachine::from_state(current_state);
364    if !sm.can_transition(ProcessingJobTransition::Fail) {
365        return (StatusCode::BAD_REQUEST, Json(ApiResponse::<ProcessingJobResponseDto>::error("Transition not allowed from current state")));
366    }
367
368    // Apply transition via partial update
369    let mut fields: HashMap<String, serde_json::Value> = HashMap::new();
370    fields.insert("status".to_string(), serde_json::Value::String("Failed".to_string()));
371
372    match service.partial_update(&id, fields).await {
373        Ok(Some(updated)) => {
374            let response: ProcessingJobResponseDto = updated.into();
375            (StatusCode::OK, Json(ApiResponse::ok(response)))
376        }
377        Ok(None) => (StatusCode::NOT_FOUND, Json(ApiResponse::<ProcessingJobResponseDto>::not_found("ProcessingJob", &id))),
378        Err(e) => (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<ProcessingJobResponseDto>::error(e.to_string()))),
379    }
380}
381
382/// Execute cancel_pending transition on a ProcessingJob.
383///
384/// POST /processing_jobs/:id/transitions/cancel_pending
385pub async fn cancel_pending_transition(
386    axum::extract::State(service): axum::extract::State<Arc<ProcessingJobService>>,
387    axum::extract::Path(id): axum::extract::Path<String>,
388    #[cfg(feature = "auth")] axum::Extension(auth): axum::Extension<AuthContext>,
389) -> impl axum::response::IntoResponse {
390    use axum::{http::StatusCode, Json};
391
392    // Get current entity
393    let entity = match service.get_by_id(&id).await {
394        Ok(Some(e)) => e,
395        Ok(None) => return (StatusCode::NOT_FOUND, Json(ApiResponse::<ProcessingJobResponseDto>::not_found("ProcessingJob", &id))),
396        Err(e) => return (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<ProcessingJobResponseDto>::error(e.to_string()))),
397    };
398
399    // Check permission (if auth enabled)
400    #[cfg(feature = "auth")]
401    {
402        let allowed_roles = ProcessingJobTransition::CancelPending.allowed_roles();
403        let has_specific_perm = auth.permissions.iter().any(|p| p == "processing_job:transition:cancel_pending");
404        let has_update_perm = auth.permissions.iter().any(|p| p == "processing_job:update");
405        let has_role = auth.roles.iter().any(|r| allowed_roles.contains(&r.as_str()));
406        if !has_specific_perm && !has_update_perm && !has_role {
407            return (StatusCode::FORBIDDEN, Json(ApiResponse::<ProcessingJobResponseDto>::error("Insufficient permissions for cancel_pending transition")));
408        }
409    }
410
411    // Create state machine from entity's actual status and validate transition
412    let current_state: ProcessingJobState = entity.status.to_string().parse()
413        .unwrap_or(ProcessingJobState::default());
414    let sm = ProcessingJobStateMachine::from_state(current_state);
415    if !sm.can_transition(ProcessingJobTransition::CancelPending) {
416        return (StatusCode::BAD_REQUEST, Json(ApiResponse::<ProcessingJobResponseDto>::error("Transition not allowed from current state")));
417    }
418
419    // Apply transition via partial update
420    let mut fields: HashMap<String, serde_json::Value> = HashMap::new();
421    fields.insert("status".to_string(), serde_json::Value::String("Cancelled".to_string()));
422
423    match service.partial_update(&id, fields).await {
424        Ok(Some(updated)) => {
425            let response: ProcessingJobResponseDto = updated.into();
426            (StatusCode::OK, Json(ApiResponse::ok(response)))
427        }
428        Ok(None) => (StatusCode::NOT_FOUND, Json(ApiResponse::<ProcessingJobResponseDto>::not_found("ProcessingJob", &id))),
429        Err(e) => (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<ProcessingJobResponseDto>::error(e.to_string()))),
430    }
431}
432
433/// Execute cancel_running transition on a ProcessingJob.
434///
435/// POST /processing_jobs/:id/transitions/cancel_running
436pub async fn cancel_running_transition(
437    axum::extract::State(service): axum::extract::State<Arc<ProcessingJobService>>,
438    axum::extract::Path(id): axum::extract::Path<String>,
439    #[cfg(feature = "auth")] axum::Extension(auth): axum::Extension<AuthContext>,
440) -> impl axum::response::IntoResponse {
441    use axum::{http::StatusCode, Json};
442
443    // Get current entity
444    let entity = match service.get_by_id(&id).await {
445        Ok(Some(e)) => e,
446        Ok(None) => return (StatusCode::NOT_FOUND, Json(ApiResponse::<ProcessingJobResponseDto>::not_found("ProcessingJob", &id))),
447        Err(e) => return (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<ProcessingJobResponseDto>::error(e.to_string()))),
448    };
449
450    // Check permission (if auth enabled)
451    #[cfg(feature = "auth")]
452    {
453        let allowed_roles = ProcessingJobTransition::CancelRunning.allowed_roles();
454        let has_specific_perm = auth.permissions.iter().any(|p| p == "processing_job:transition:cancel_running");
455        let has_update_perm = auth.permissions.iter().any(|p| p == "processing_job:update");
456        let has_role = auth.roles.iter().any(|r| allowed_roles.contains(&r.as_str()));
457        if !has_specific_perm && !has_update_perm && !has_role {
458            return (StatusCode::FORBIDDEN, Json(ApiResponse::<ProcessingJobResponseDto>::error("Insufficient permissions for cancel_running transition")));
459        }
460    }
461
462    // Create state machine from entity's actual status and validate transition
463    let current_state: ProcessingJobState = entity.status.to_string().parse()
464        .unwrap_or(ProcessingJobState::default());
465    let sm = ProcessingJobStateMachine::from_state(current_state);
466    if !sm.can_transition(ProcessingJobTransition::CancelRunning) {
467        return (StatusCode::BAD_REQUEST, Json(ApiResponse::<ProcessingJobResponseDto>::error("Transition not allowed from current state")));
468    }
469
470    // Apply transition via partial update
471    let mut fields: HashMap<String, serde_json::Value> = HashMap::new();
472    fields.insert("status".to_string(), serde_json::Value::String("Cancelled".to_string()));
473
474    match service.partial_update(&id, fields).await {
475        Ok(Some(updated)) => {
476            let response: ProcessingJobResponseDto = updated.into();
477            (StatusCode::OK, Json(ApiResponse::ok(response)))
478        }
479        Ok(None) => (StatusCode::NOT_FOUND, Json(ApiResponse::<ProcessingJobResponseDto>::not_found("ProcessingJob", &id))),
480        Err(e) => (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<ProcessingJobResponseDto>::error(e.to_string()))),
481    }
482}
483
484/// Execute retry transition on a ProcessingJob.
485///
486/// POST /processing_jobs/:id/transitions/retry
487pub async fn retry_transition(
488    axum::extract::State(service): axum::extract::State<Arc<ProcessingJobService>>,
489    axum::extract::Path(id): axum::extract::Path<String>,
490    #[cfg(feature = "auth")] axum::Extension(auth): axum::Extension<AuthContext>,
491) -> impl axum::response::IntoResponse {
492    use axum::{http::StatusCode, Json};
493
494    // Get current entity
495    let entity = match service.get_by_id(&id).await {
496        Ok(Some(e)) => e,
497        Ok(None) => return (StatusCode::NOT_FOUND, Json(ApiResponse::<ProcessingJobResponseDto>::not_found("ProcessingJob", &id))),
498        Err(e) => return (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<ProcessingJobResponseDto>::error(e.to_string()))),
499    };
500
501    // Check permission (if auth enabled)
502    #[cfg(feature = "auth")]
503    {
504        let allowed_roles = ProcessingJobTransition::Retry.allowed_roles();
505        let has_specific_perm = auth.permissions.iter().any(|p| p == "processing_job:transition:retry");
506        let has_update_perm = auth.permissions.iter().any(|p| p == "processing_job:update");
507        let has_role = auth.roles.iter().any(|r| allowed_roles.contains(&r.as_str()));
508        if !has_specific_perm && !has_update_perm && !has_role {
509            return (StatusCode::FORBIDDEN, Json(ApiResponse::<ProcessingJobResponseDto>::error("Insufficient permissions for retry transition")));
510        }
511    }
512
513    // Create state machine from entity's actual status and validate transition
514    let current_state: ProcessingJobState = entity.status.to_string().parse()
515        .unwrap_or(ProcessingJobState::default());
516    let sm = ProcessingJobStateMachine::from_state(current_state);
517    if !sm.can_transition(ProcessingJobTransition::Retry) {
518        return (StatusCode::BAD_REQUEST, Json(ApiResponse::<ProcessingJobResponseDto>::error("Transition not allowed from current state")));
519    }
520
521    // Apply transition via partial update
522    let mut fields: HashMap<String, serde_json::Value> = HashMap::new();
523    fields.insert("status".to_string(), serde_json::Value::String("Pending".to_string()));
524
525    match service.partial_update(&id, fields).await {
526        Ok(Some(updated)) => {
527            let response: ProcessingJobResponseDto = updated.into();
528            (StatusCode::OK, Json(ApiResponse::ok(response)))
529        }
530        Ok(None) => (StatusCode::NOT_FOUND, Json(ApiResponse::<ProcessingJobResponseDto>::not_found("ProcessingJob", &id))),
531        Err(e) => (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<ProcessingJobResponseDto>::error(e.to_string()))),
532    }
533}
534
535/// Execute reprocess transition on a ProcessingJob.
536///
537/// POST /processing_jobs/:id/transitions/reprocess
538pub async fn reprocess_transition(
539    axum::extract::State(service): axum::extract::State<Arc<ProcessingJobService>>,
540    axum::extract::Path(id): axum::extract::Path<String>,
541    #[cfg(feature = "auth")] axum::Extension(auth): axum::Extension<AuthContext>,
542) -> impl axum::response::IntoResponse {
543    use axum::{http::StatusCode, Json};
544
545    // Get current entity
546    let entity = match service.get_by_id(&id).await {
547        Ok(Some(e)) => e,
548        Ok(None) => return (StatusCode::NOT_FOUND, Json(ApiResponse::<ProcessingJobResponseDto>::not_found("ProcessingJob", &id))),
549        Err(e) => return (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<ProcessingJobResponseDto>::error(e.to_string()))),
550    };
551
552    // Check permission (if auth enabled)
553    #[cfg(feature = "auth")]
554    {
555        let allowed_roles = ProcessingJobTransition::Reprocess.allowed_roles();
556        let has_specific_perm = auth.permissions.iter().any(|p| p == "processing_job:transition:reprocess");
557        let has_update_perm = auth.permissions.iter().any(|p| p == "processing_job:update");
558        let has_role = auth.roles.iter().any(|r| allowed_roles.contains(&r.as_str()));
559        if !has_specific_perm && !has_update_perm && !has_role {
560            return (StatusCode::FORBIDDEN, Json(ApiResponse::<ProcessingJobResponseDto>::error("Insufficient permissions for reprocess transition")));
561        }
562    }
563
564    // Create state machine from entity's actual status and validate transition
565    let current_state: ProcessingJobState = entity.status.to_string().parse()
566        .unwrap_or(ProcessingJobState::default());
567    let sm = ProcessingJobStateMachine::from_state(current_state);
568    if !sm.can_transition(ProcessingJobTransition::Reprocess) {
569        return (StatusCode::BAD_REQUEST, Json(ApiResponse::<ProcessingJobResponseDto>::error("Transition not allowed from current state")));
570    }
571
572    // Apply transition via partial update
573    let mut fields: HashMap<String, serde_json::Value> = HashMap::new();
574    fields.insert("status".to_string(), serde_json::Value::String("Pending".to_string()));
575
576    match service.partial_update(&id, fields).await {
577        Ok(Some(updated)) => {
578            let response: ProcessingJobResponseDto = updated.into();
579            (StatusCode::OK, Json(ApiResponse::ok(response)))
580        }
581        Ok(None) => (StatusCode::NOT_FOUND, Json(ApiResponse::<ProcessingJobResponseDto>::not_found("ProcessingJob", &id))),
582        Err(e) => (StatusCode::INTERNAL_SERVER_ERROR, Json(ApiResponse::<ProcessingJobResponseDto>::error(e.to_string()))),
583    }
584}
585
586/// Create routes for state transitions.
587pub fn create_processing_job_transition_routes(service: Arc<ProcessingJobService>) -> Router {
588    use axum::routing::post;
589
590    Router::new()
591        .route("/processing_jobs/:id/transitions/start", post(start_transition))
592        .route("/processing_jobs/:id/transitions/complete", post(complete_transition))
593        .route("/processing_jobs/:id/transitions/fail", post(fail_transition))
594        .route("/processing_jobs/:id/transitions/cancel_pending", post(cancel_pending_transition))
595        .route("/processing_jobs/:id/transitions/cancel_running", post(cancel_running_transition))
596        .route("/processing_jobs/:id/transitions/retry", post(retry_transition))
597        .route("/processing_jobs/:id/transitions/reprocess", post(reprocess_transition))
598        .with_state(service)
599}