oxedyne/fe2o3/fe2o3_datime/src/schedule/task.rs
11.5 KiB, 54 runs
created by r1870400018:8498, 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 definitions for the scheduling system. |
| 2 | //! |
| 3 | //! A task pairs an action with a time to run it, a priority, a retry policy |
| 4 | //! and, for recurring tasks, a pattern. Each run is recorded, and the recent |
| 5 | //! records are what the success rate and average duration are drawn from. |
| 6 | //! |
| 7 | //! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\ |
| 8 | //! Anthropic Claude |
| 9 | |
| 10 | use oxedyne_fe2o3_core::prelude::*; |
| 11 | use crate::{ |
| 12 | time::{CalClock, CalClockZone}, |
| 13 | schedule::{RecurrencePattern, Action}, |
| 14 | }; |
| 15 | use std::{ |
| 16 | fmt, |
| 17 | sync::atomic::{AtomicU64, Ordering}, |
| 18 | }; |
| 19 | |
| 20 | #[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] |
| 21 | pub struct TaskId(u64); |
| 22 | |
| 23 | impl TaskId { |
| 24 | fn next() -> Self { |
| 25 | static COUNTER: AtomicU64 = AtomicU64::new(1); |
| 26 | TaskId(COUNTER.fetch_add(1, Ordering::Relaxed)) |
| 27 | } |
| 28 | } |
| 29 | |
| 30 | impl fmt::Display for TaskId { |
| 31 | fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { |
| 32 | write!(f, "task-{}", self.0) |
| 33 | } |
| 34 | } |
| 35 | |
| 36 | #[derive(Debug, Clone, PartialEq, Eq)] |
| 37 | pub enum TaskStatus { |
| 38 | Pending, |
| 39 | Running, |
| 40 | Completed, |
| 41 | Failed(String), // reason |
| 42 | Cancelled, // never started |
| 43 | Skipped, // its turn came and went |
| 44 | } |
| 45 | |
| 46 | #[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)] |
| 47 | pub enum TaskPriority { |
| 48 | Low = 1, |
| 49 | Normal = 2, |
| 50 | High = 3, |
| 51 | Critical = 4, |
| 52 | } |
| 53 | |
| 54 | impl Default for TaskPriority { |
| 55 | fn default() -> Self { |
| 56 | TaskPriority::Normal |
| 57 | } |
| 58 | } |
| 59 | |
| 60 | #[derive(Debug, Clone)] |
| 61 | pub struct TaskConfig { |
| 62 | pub timeout_millis: Option<u64>, |
| 63 | pub max_retries: u32, |
| 64 | pub retry_delay_millis: u64, |
| 65 | pub continue_on_failure: bool, // keep scheduling after a failure |
| 66 | pub allow_concurrent: bool, // more than one instance at a time |
| 67 | } |
| 68 | |
| 69 | impl Default for TaskConfig { |
| 70 | fn default() -> Self { |
| 71 | TaskConfig { |
| 72 | timeout_millis: Some(300_000), // 5 minutes default timeout |
| 73 | max_retries: 3, |
| 74 | retry_delay_millis: 1000, // 1 second delay |
| 75 | continue_on_failure: true, |
| 76 | allow_concurrent: false, |
| 77 | } |
| 78 | } |
| 79 | } |
| 80 | |
| 81 | #[derive(Debug)] |
| 82 | pub struct Task { |
| 83 | pub id: TaskId, |
| 84 | pub name: String, |
| 85 | pub description: Option<String>, |
| 86 | pub zone: CalClockZone, |
| 87 | pub scheduled_time: CalClock, |
| 88 | pub recurrence: Option<RecurrencePattern>, // None means one-time |
| 89 | pub priority: TaskPriority, |
| 90 | pub config: TaskConfig, |
| 91 | pub status: TaskStatus, |
| 92 | pub action: Box<dyn Action>, |
| 93 | pub execution_count: u32, |
| 94 | pub last_execution: Option<CalClock>, |
| 95 | pub next_execution: Option<CalClock>, |
| 96 | pub execution_history: Vec<TaskExecution>, // recent only |
| 97 | } |
| 98 | |
| 99 | impl Clone for Task { |
| 100 | fn clone(&self) -> Self { |
| 101 | Task { |
| 102 | id: self.id, |
| 103 | name: self.name.clone(), |
| 104 | description: self.description.clone(), |
| 105 | zone: self.zone.clone(), |
| 106 | scheduled_time: self.scheduled_time.clone(), |
| 107 | recurrence: self.recurrence.clone(), |
| 108 | priority: self.priority, |
| 109 | config: self.config.clone(), |
| 110 | status: self.status.clone(), |
| 111 | action: self.action.box_clone(), |
| 112 | execution_count: self.execution_count, |
| 113 | last_execution: self.last_execution.clone(), |
| 114 | next_execution: self.next_execution.clone(), |
| 115 | execution_history: self.execution_history.clone(), |
| 116 | } |
| 117 | } |
| 118 | } |
| 119 | |
| 120 | #[derive(Debug, Clone)] |
| 121 | pub struct TaskExecution { |
| 122 | pub started_at: CalClock, |
| 123 | pub completed_at: Option<CalClock>, // None while still running |
| 124 | pub duration_millis: Option<u64>, |
| 125 | pub status: TaskStatus, |
| 126 | pub error_message: Option<String>, |
| 127 | pub retry_count: u32, |
| 128 | } |
| 129 | |
| 130 | impl Task { |
| 131 | pub fn new<S: Into<String>>(name: S, zone: CalClockZone) -> TaskBuilder { |
| 132 | TaskBuilder::new(name.into(), zone) |
| 133 | } |
| 134 | |
| 135 | pub fn is_ready_to_execute(&self, current_time: &CalClock) -> bool { |
| 136 | match self.status { |
| 137 | TaskStatus::Pending => current_time >= &self.scheduled_time, |
| 138 | _ => false, |
| 139 | } |
| 140 | } |
| 141 | |
| 142 | pub fn calculate_next_execution(&self) -> Outcome<Option<CalClock>> { |
| 143 | if let Some(ref pattern) = self.recurrence { |
| 144 | pattern.next_execution(&self.scheduled_time, &self.zone) |
| 145 | } else { |
| 146 | Ok(None) |
| 147 | } |
| 148 | } |
| 149 | |
| 150 | pub fn advance_to_next_execution(&mut self) -> Outcome<()> { |
| 151 | if let Some(next_time) = res!(self.calculate_next_execution()) { |
| 152 | self.scheduled_time = next_time; |
| 153 | self.next_execution = res!(self.calculate_next_execution()); |
| 154 | self.status = TaskStatus::Pending; |
| 155 | } |
| 156 | Ok(()) |
| 157 | } |
| 158 | |
| 159 | pub fn record_execution(&mut self, execution: TaskExecution) { |
| 160 | const MAX_HISTORY: usize = 10; // Keep last 10 executions |
| 161 | |
| 162 | self.execution_count += 1; |
| 163 | self.last_execution = Some(execution.started_at.clone()); |
| 164 | self.execution_history.push(execution); |
| 165 | |
| 166 | // Keep only recent executions |
| 167 | if self.execution_history.len() > MAX_HISTORY { |
| 168 | self.execution_history.drain(0..self.execution_history.len() - MAX_HISTORY); |
| 169 | } |
| 170 | } |
| 171 | |
| 172 | /// Over the retained history only. |
| 173 | pub fn average_execution_duration(&self) -> Option<u64> { |
| 174 | let completed: Vec<_> = self.execution_history.iter() |
| 175 | .filter_map(|exec| exec.duration_millis) |
| 176 | .collect(); |
| 177 | |
| 178 | if completed.is_empty() { |
| 179 | None |
| 180 | } else { |
| 181 | Some(completed.iter().sum::<u64>() / completed.len() as u64) |
| 182 | } |
| 183 | } |
| 184 | |
| 185 | /// A percentage, over the retained history only. |
| 186 | pub fn success_rate(&self) -> f64 { |
| 187 | if self.execution_history.is_empty() { |
| 188 | 0.0 |
| 189 | } else { |
| 190 | let successful = self.execution_history.iter() |
| 191 | .filter(|exec| matches!(exec.status, TaskStatus::Completed)) |
| 192 | .count(); |
| 193 | (successful as f64 / self.execution_history.len() as f64) * 100.0 |
| 194 | } |
| 195 | } |
| 196 | } |
| 197 | |
| 198 | pub struct TaskBuilder { |
| 199 | name: String, |
| 200 | description: Option<String>, |
| 201 | zone: CalClockZone, |
| 202 | scheduled_time: Option<CalClock>, |
| 203 | recurrence: Option<RecurrencePattern>, |
| 204 | priority: TaskPriority, |
| 205 | config: TaskConfig, |
| 206 | action: Option<Box<dyn Action>>, |
| 207 | } |
| 208 | |
| 209 | impl TaskBuilder { |
| 210 | pub fn new(name: String, zone: CalClockZone) -> Self { |
| 211 | TaskBuilder { |
| 212 | name, |
| 213 | description: None, |
| 214 | zone, |
| 215 | scheduled_time: None, |
| 216 | recurrence: None, |
| 217 | priority: TaskPriority::default(), |
| 218 | config: TaskConfig::default(), |
| 219 | action: None, |
| 220 | } |
| 221 | } |
| 222 | |
| 223 | pub fn description<S: Into<String>>(mut self, desc: S) -> Self { |
| 224 | self.description = Some(desc.into()); |
| 225 | self |
| 226 | } |
| 227 | |
| 228 | pub fn at_time(mut self, hour: u8, minute: u8, second: u8) -> Self { |
| 229 | if let Ok(time) = CalClock::new(2024, 1, 1, hour, minute, second, 0, self.zone.clone()) { |
| 230 | self.scheduled_time = Some(time); |
| 231 | } |
| 232 | self |
| 233 | } |
| 234 | |
| 235 | pub fn on_date(mut self, year: i32, month: u8, day: u8) -> Self { |
| 236 | if let Some(existing_time) = self.scheduled_time.take() { |
| 237 | if let Ok(time) = CalClock::new( |
| 238 | year, month, day, |
| 239 | existing_time.hour(), |
| 240 | existing_time.minute(), |
| 241 | existing_time.second(), |
| 242 | existing_time.nanosecond(), |
| 243 | self.zone.clone() |
| 244 | ) { |
| 245 | self.scheduled_time = Some(time); |
| 246 | } |
| 247 | } else if let Ok(time) = CalClock::new(year, month, day, 0, 0, 0, 0, self.zone.clone()) { |
| 248 | self.scheduled_time = Some(time); |
| 249 | } |
| 250 | self |
| 251 | } |
| 252 | |
| 253 | pub fn at(mut self, time: CalClock) -> Self { |
| 254 | self.scheduled_time = Some(time); |
| 255 | self |
| 256 | } |
| 257 | |
| 258 | pub fn recurrence(mut self, pattern: RecurrencePattern) -> Self { |
| 259 | self.recurrence = Some(pattern); |
| 260 | self |
| 261 | } |
| 262 | |
| 263 | pub fn priority(mut self, priority: TaskPriority) -> Self { |
| 264 | self.priority = priority; |
| 265 | self |
| 266 | } |
| 267 | |
| 268 | pub fn config(mut self, config: TaskConfig) -> Self { |
| 269 | self.config = config; |
| 270 | self |
| 271 | } |
| 272 | |
| 273 | pub fn with_action<A: Action + 'static>(mut self, action: A) -> Self { |
| 274 | self.action = Some(Box::new(action)); |
| 275 | self |
| 276 | } |
| 277 | |
| 278 | pub fn build(self) -> Outcome<Task> { |
| 279 | let scheduled_time = res!(self.scheduled_time |
| 280 | .ok_or_else(|| err!("Task scheduled time is required"; Invalid, Input))); |
| 281 | |
| 282 | let action = res!(self.action |
| 283 | .ok_or_else(|| err!("Task action is required"; Invalid, Input))); |
| 284 | |
| 285 | let next_execution = if self.recurrence.is_some() { |
| 286 | // For recurring tasks, calculate the first next execution |
| 287 | if let Some(ref pattern) = self.recurrence { |
| 288 | res!(pattern.next_execution(&scheduled_time, &self.zone)) |
| 289 | } else { |
| 290 | None |
| 291 | } |
| 292 | } else { |
| 293 | None |
| 294 | }; |
| 295 | |
| 296 | Ok(Task { |
| 297 | id: TaskId::next(), |
| 298 | name: self.name, |
| 299 | description: self.description, |
| 300 | zone: self.zone, |
| 301 | scheduled_time, |
| 302 | recurrence: self.recurrence, |
| 303 | priority: self.priority, |
| 304 | config: self.config, |
| 305 | status: TaskStatus::Pending, |
| 306 | action, |
| 307 | execution_count: 0, |
| 308 | last_execution: None, |
| 309 | next_execution, |
| 310 | execution_history: Vec::new(), |
| 311 | }) |
| 312 | } |
| 313 | } |
| 314 | |
| 315 | #[cfg(test)] |
| 316 | mod tests { |
| 317 | use super::*; |
| 318 | use crate::schedule::action::CallbackAction; |
| 319 | |
| 320 | #[test] |
| 321 | fn test_task_builder() { |
| 322 | let zone = CalClockZone::utc(); |
| 323 | let callback = CallbackAction::new(|| { |
| 324 | println!("Test action executed"); |
| 325 | Ok(()) |
| 326 | }); |
| 327 | |
| 328 | let task = Task::new("test_task", zone.clone()) |
| 329 | .description("A test task") |
| 330 | .at_time(14, 30, 0) |
| 331 | .on_date(2024, 6, 15) |
| 332 | .priority(TaskPriority::High) |
| 333 | .with_action(callback) |
| 334 | .build() |
| 335 | .expect("Failed to build task"); |
| 336 | |
| 337 | assert_eq!(task.name, "test_task"); |
| 338 | assert_eq!(task.priority, TaskPriority::High); |
| 339 | assert_eq!(task.status, TaskStatus::Pending); |
| 340 | assert_eq!(task.scheduled_time.year(), 2024); |
| 341 | assert_eq!(task.scheduled_time.month(), 6); |
| 342 | assert_eq!(task.scheduled_time.day(), 15); |
| 343 | assert_eq!(task.scheduled_time.hour(), 14); |
| 344 | assert_eq!(task.scheduled_time.minute(), 30); |
| 345 | } |
| 346 | |
| 347 | #[test] |
| 348 | fn test_task_id_uniqueness() { |
| 349 | let id1 = TaskId::next(); |
| 350 | let id2 = TaskId::next(); |
| 351 | assert_ne!(id1, id2); |
| 352 | } |
| 353 | |
| 354 | #[test] |
| 355 | fn test_task_execution_recording() { |
| 356 | let zone = CalClockZone::utc(); |
| 357 | let callback = CallbackAction::new(|| Ok(())); |
| 358 | |
| 359 | let mut task = Task::new("test", zone.clone()) |
| 360 | .at_time(12, 0, 0) |
| 361 | .on_date(2024, 1, 1) |
| 362 | .with_action(callback) |
| 363 | .build() |
| 364 | .unwrap(); |
| 365 | |
| 366 | let execution = TaskExecution { |
| 367 | started_at: CalClock::new(2024, 1, 1, 12, 0, 0, 0, zone.clone()).unwrap(), |
| 368 | completed_at: Some(CalClock::new(2024, 1, 1, 12, 0, 5, 0, zone).unwrap()), |
| 369 | duration_millis: Some(5000), |
| 370 | status: TaskStatus::Completed, |
| 371 | error_message: None, |
| 372 | retry_count: 0, |
| 373 | }; |
| 374 | |
| 375 | task.record_execution(execution); |
| 376 | |
| 377 | assert_eq!(task.execution_count, 1); |
| 378 | assert_eq!(task.execution_history.len(), 1); |
| 379 | assert_eq!(task.success_rate(), 100.0); |
| 380 | assert_eq!(task.average_execution_duration(), Some(5000)); |
| 381 | } |
| 382 | } |