oxedyne/fe2o3/fe2o3_datime/src/schedule/executor.rs
8.4 KiB, 40 runs
created by r1870400018:8490, which is this file's identity for as long as the history lasts, whatever it is later renamed to
download · who wrote it · its history
| 1 | //! Task execution engine for the scheduling system. |
| 2 | //! |
| 3 | //! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\ |
| 4 | //! Anthropic Claude |
| 5 | |
| 6 | use oxedyne_fe2o3_core::prelude::*; |
| 7 | use crate::{ |
| 8 | schedule::{Task, TaskStatus, TaskExecution, ActionContext, ActionResult}, |
| 9 | time::CalClock, |
| 10 | }; |
| 11 | use std::{ |
| 12 | time::{Instant, Duration}, |
| 13 | thread, |
| 14 | }; |
| 15 | |
| 16 | #[derive(Debug, Clone)] |
| 17 | pub struct ExecutionResult { |
| 18 | pub success: bool, |
| 19 | pub duration_millis: u64, |
| 20 | pub error_message: Option<String>, |
| 21 | pub action_result: Option<ActionResult>, |
| 22 | } |
| 23 | |
| 24 | #[derive(Debug, Clone, Default)] |
| 25 | pub struct ExecutionStats { |
| 26 | pub total_executions: u64, |
| 27 | pub successful_executions: u64, |
| 28 | pub failed_executions: u64, |
| 29 | pub avg_execution_time_millis: u64, |
| 30 | total_execution_time_millis: u64, // the running sum behind the average |
| 31 | } |
| 32 | |
| 33 | #[derive(Debug)] |
| 34 | pub struct TaskExecutor { |
| 35 | stats: ExecutionStats, |
| 36 | } |
| 37 | |
| 38 | impl TaskExecutor { |
| 39 | pub fn new() -> Self { |
| 40 | TaskExecutor { |
| 41 | stats: ExecutionStats::default(), |
| 42 | } |
| 43 | } |
| 44 | |
| 45 | pub fn execute(&mut self, task: &Task) -> Outcome<ExecutionResult> { |
| 46 | let start_time = Instant::now(); |
| 47 | let execution_time = res!(CalClock::now_utc()); |
| 48 | |
| 49 | // Create action context |
| 50 | let context = ActionContext { |
| 51 | scheduled_time: task.scheduled_time.clone(), |
| 52 | execution_time: execution_time.clone(), |
| 53 | task_name: task.name.clone(), |
| 54 | execution_count: task.execution_count, |
| 55 | is_retry: false, // TODO: Track retry state properly |
| 56 | retry_count: 0, |
| 57 | }; |
| 58 | |
| 59 | // Update statistics |
| 60 | self.stats.total_executions += 1; |
| 61 | |
| 62 | // Prepare the action |
| 63 | if let Err(e) = task.action.prepare(&context) { |
| 64 | let error_msg = format!("Action preparation failed: {}", e); |
| 65 | self.stats.failed_executions += 1; |
| 66 | return Ok(ExecutionResult { |
| 67 | success: false, |
| 68 | duration_millis: start_time.elapsed().as_millis() as u64, |
| 69 | error_message: Some(error_msg), |
| 70 | action_result: None, |
| 71 | }); |
| 72 | } |
| 73 | |
| 74 | // Execute the action with timeout handling |
| 75 | let action_result = if let Some(timeout_millis) = task.config.timeout_millis { |
| 76 | self.execute_with_timeout(&*task.action, &context, timeout_millis) |
| 77 | } else { |
| 78 | task.action.execute(&context) |
| 79 | }; |
| 80 | |
| 81 | let duration = start_time.elapsed(); |
| 82 | let duration_millis = duration.as_millis() as u64; |
| 83 | |
| 84 | // Update timing statistics |
| 85 | self.stats.total_execution_time_millis += duration_millis; |
| 86 | self.stats.avg_execution_time_millis = |
| 87 | self.stats.total_execution_time_millis / self.stats.total_executions; |
| 88 | |
| 89 | // Process action result |
| 90 | let (success, error_message) = match &action_result { |
| 91 | Ok(ActionResult::Success) => { |
| 92 | self.stats.successful_executions += 1; |
| 93 | (true, None) |
| 94 | }, |
| 95 | Ok(ActionResult::Warning(msg)) => { |
| 96 | self.stats.successful_executions += 1; |
| 97 | (true, Some(format!("Warning: {}", msg))) |
| 98 | }, |
| 99 | Ok(ActionResult::Error(action_error)) => { |
| 100 | self.stats.failed_executions += 1; |
| 101 | (false, Some(format!("Action error: {}", action_error))) |
| 102 | }, |
| 103 | Err(e) => { |
| 104 | self.stats.failed_executions += 1; |
| 105 | (false, Some(format!("Execution failed: {}", e))) |
| 106 | } |
| 107 | }; |
| 108 | |
| 109 | // Cleanup the action |
| 110 | if let Ok(ref result) = action_result { |
| 111 | let _ = task.action.cleanup(&context, result); |
| 112 | } |
| 113 | |
| 114 | Ok(ExecutionResult { |
| 115 | success, |
| 116 | duration_millis, |
| 117 | error_message, |
| 118 | action_result: action_result.ok(), |
| 119 | }) |
| 120 | } |
| 121 | |
| 122 | /// The action is run to completion, so a timeout is reported after the |
| 123 | /// fact rather than cutting the action short. |
| 124 | fn execute_with_timeout( |
| 125 | &self, |
| 126 | action: &dyn crate::schedule::Action, |
| 127 | context: &ActionContext, |
| 128 | timeout_millis: u64 |
| 129 | ) -> Outcome<ActionResult> { |
| 130 | use std::sync::mpsc; |
| 131 | |
| 132 | let (_tx, _rx) = mpsc::channel::<Outcome<()>>(); |
| 133 | let action_context = context.clone(); |
| 134 | |
| 135 | // We need to work around the fact that Action is not Clone |
| 136 | // For now, we'll execute directly and add timeout simulation |
| 137 | let start = Instant::now(); |
| 138 | let result = action.execute(&action_context); |
| 139 | let elapsed = start.elapsed(); |
| 140 | |
| 141 | if elapsed.as_millis() as u64 > timeout_millis { |
| 142 | Err(err!("Action execution timed out after {}ms", timeout_millis; Timeout, Input)) |
| 143 | } else { |
| 144 | result |
| 145 | } |
| 146 | } |
| 147 | |
| 148 | pub fn execute_with_retry(&mut self, task: &mut Task) -> Outcome<ExecutionResult> { |
| 149 | let mut last_result = None; |
| 150 | let max_retries = task.config.max_retries; |
| 151 | |
| 152 | for attempt in 0..=max_retries { |
| 153 | let is_retry = attempt > 0; |
| 154 | |
| 155 | if is_retry { |
| 156 | // Wait before retry |
| 157 | thread::sleep(Duration::from_millis(task.config.retry_delay_millis)); |
| 158 | } |
| 159 | |
| 160 | // Update task status |
| 161 | task.status = TaskStatus::Running; |
| 162 | |
| 163 | // Execute the task |
| 164 | let result = self.execute(task); |
| 165 | |
| 166 | match &result { |
| 167 | Ok(exec_result) if exec_result.success => { |
| 168 | // Success - record execution and return |
| 169 | task.status = TaskStatus::Completed; |
| 170 | self.record_task_execution(task, exec_result, attempt); |
| 171 | return result; |
| 172 | }, |
| 173 | Ok(exec_result) => { |
| 174 | // Failed execution |
| 175 | if attempt == max_retries { |
| 176 | // Final attempt failed |
| 177 | task.status = TaskStatus::Failed( |
| 178 | exec_result.error_message.clone() |
| 179 | .unwrap_or_else(|| "Unknown error".to_string()) |
| 180 | ); |
| 181 | self.record_task_execution(task, exec_result, attempt); |
| 182 | return result; |
| 183 | } else { |
| 184 | // Will retry |
| 185 | last_result = Some(exec_result.clone()); |
| 186 | } |
| 187 | }, |
| 188 | Err(_) => { |
| 189 | // Critical error - don't retry |
| 190 | task.status = TaskStatus::Failed("Critical execution error".to_string()); |
| 191 | return result; |
| 192 | } |
| 193 | } |
| 194 | } |
| 195 | |
| 196 | // This should never be reached due to the logic above |
| 197 | Ok(last_result.unwrap_or_else(|| ExecutionResult { |
| 198 | success: false, |
| 199 | duration_millis: 0, |
| 200 | error_message: Some("Unexpected execution state".to_string()), |
| 201 | action_result: None, |
| 202 | })) |
| 203 | } |
| 204 | |
| 205 | fn record_task_execution(&self, task: &mut Task, result: &ExecutionResult, retry_count: u32) { |
| 206 | let started_at = CalClock::now_utc() |
| 207 | .or_else(|_| CalClock::new(2024, 1, 1, 0, 0, 0, 0, crate::time::CalClockZone::utc())) |
| 208 | .unwrap_or_else(|_| task.scheduled_time.clone()); |
| 209 | let completed_at = CalClock::now_utc() |
| 210 | .or_else(|_| CalClock::new(2024, 1, 1, 0, 0, 0, 0, crate::time::CalClockZone::utc())) |
| 211 | .unwrap_or_else(|_| task.scheduled_time.clone()); |
| 212 | |
| 213 | |
| 214 | let execution = TaskExecution { |
| 215 | started_at, |
| 216 | completed_at: Some(completed_at), |
| 217 | duration_millis: Some(result.duration_millis), |
| 218 | status: if result.success { |
| 219 | TaskStatus::Completed |
| 220 | } else { |
| 221 | TaskStatus::Failed( |
| 222 | result.error_message.clone() |
| 223 | .unwrap_or_else(|| "Unknown error".to_string()) |
| 224 | ) |
| 225 | }, |
| 226 | error_message: result.error_message.clone(), |
| 227 | retry_count, |
| 228 | }; |
| 229 | |
| 230 | task.record_execution(execution); |
| 231 | } |
| 232 | |
| 233 | pub fn stats(&self) -> &ExecutionStats { |
| 234 | &self.stats |
| 235 | } |
| 236 | |
| 237 | pub fn reset_stats(&mut self) { |
| 238 | self.stats = ExecutionStats::default(); |
| 239 | } |
| 240 | |
| 241 | /// A percentage. |
| 242 | pub fn success_rate(&self) -> f64 { |
| 243 | if self.stats.total_executions == 0 { |
| 244 | 0.0 |
| 245 | } else { |
| 246 | (self.stats.successful_executions as f64 / self.stats.total_executions as f64) * 100.0 |
| 247 | } |
| 248 | } |
| 249 | } |