Background loop that detects and fails orphaned tasks. Tasks are considered orphaned if: - Status is 'assigned' or 'active' AND - The rover is not currently connected AND - The task has been in that state for > TASK_TIMEOUT_MINUTES
(state: SharedState)
| 600 | /// - The rover is not currently connected AND |
| 601 | /// - The task has been in that state for > TASK_TIMEOUT_MINUTES |
| 602 | async fn task_cleanup_loop(state: SharedState) { |
| 603 | use std::time::Duration; |
| 604 | use tokio::time::interval; |
| 605 | |
| 606 | // Check every minute |
| 607 | let mut ticker = interval(Duration::from_secs(60)); |
| 608 | // Configurable timeout (default 5 minutes) |
| 609 | let timeout_minutes: i64 = std::env::var("TASK_TIMEOUT_MINUTES") |
| 610 | .ok() |
| 611 | .and_then(|v| v.parse().ok()) |
| 612 | .unwrap_or(5); |
| 613 | |
| 614 | info!(timeout_minutes, "Starting task cleanup loop"); |
| 615 | |
| 616 | loop { |
| 617 | ticker.tick().await; |
| 618 | |
| 619 | // Get list of currently connected rover IDs |
| 620 | let connected_rovers: Vec<String> = { |
| 621 | let rovers = state.rovers.read().await; |
| 622 | rovers.keys().cloned().collect() |
| 623 | }; |
| 624 | |
| 625 | // Find orphaned tasks: active/assigned tasks for rovers that are NOT connected |
| 626 | // and have been stuck for longer than the timeout |
| 627 | let orphaned_tasks: Vec<Task> = sqlx::query_as( |
| 628 | r#" |
| 629 | UPDATE tasks |
| 630 | SET status = 'failed', error = 'Task timeout - rover not connected', ended_at = now() |
| 631 | WHERE status IN ('assigned', 'active') |
| 632 | AND rover_id != ALL($1) |
| 633 | AND ( |
| 634 | (started_at IS NOT NULL AND started_at < now() - make_interval(mins => $2)) |
| 635 | OR (started_at IS NULL AND created_at < now() - make_interval(mins => $2)) |
| 636 | ) |
| 637 | RETURNING id, mission_id, rover_id, status, progress, waypoint, lap, error, created_at, started_at, ended_at |
| 638 | "#, |
| 639 | ) |
| 640 | .bind(&connected_rovers) |
| 641 | .bind(timeout_minutes) |
| 642 | .fetch_all(&state.db) |
| 643 | .await |
| 644 | .unwrap_or_default(); |
| 645 | |
| 646 | for task in orphaned_tasks { |
| 647 | warn!( |
| 648 | task_id = %task.id, |
| 649 | rover_id = %task.rover_id, |
| 650 | "Marked orphaned task as failed" |
| 651 | ); |
| 652 | state.broadcast(BroadcastMessage::TaskUpdate { task }); |
| 653 | } |
| 654 | } |
| 655 | } |
| 656 | |
| 657 | // ============================================================================= |
| 658 | // Health Check |