MCPcopy Create free account
hub / github.com/ecto/muni / task_cleanup_loop

Function task_cleanup_loop

depot/dispatch/src/main.rs:602–655  ·  view source on GitHub ↗

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)

Source from the content-addressed store, hash-verified

600/// - The rover is not currently connected AND
601/// - The task has been in that state for > TASK_TIMEOUT_MINUTES
602async 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

Callers 1

mainFunction · 0.85

Calls 4

okMethod · 0.80
broadcastMethod · 0.80
parseMethod · 0.45
tickMethod · 0.45

Tested by

no test coverage detected