Oregami
Repositories/oxedyne/fe2o3

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
6use oxedyne_fe2o3_core::prelude::*;
7use crate::{
8 schedule::{Task, TaskStatus, TaskExecution, ActionContext, ActionResult},
9 time::CalClock,
10};
11use std::{
12 time::{Instant, Duration},
13 thread,
14};
15
16#[derive(Debug, Clone)]
17pub 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)]
25pub 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)]
34pub struct TaskExecutor {
35 stats: ExecutionStats,
36}
37
38impl 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}