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::{ProcessingJobService, ServiceError};
27
28use crate::presentation::dto::{CreateProcessingJobDto, UpdateProcessingJobDto, PatchProcessingJobDto, ProcessingJobResponseDto};
30
31use crate::domain::state_machine::{ProcessingJobState, ProcessingJobStateMachine, ProcessingJobTransition};
32
33#[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 #[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
103pub 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
114pub 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
153pub 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
165pub 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#[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
225pub 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 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 #[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 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 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
280pub 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 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 #[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 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 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
331pub 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 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 #[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 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 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
382pub 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 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 #[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 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 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
433pub 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 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 #[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 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 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
484pub 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 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 #[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 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 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
535pub 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 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 #[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 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 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
586pub 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}