Oregami
Repositories/oxedyne/fe2o3

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
10use oxedyne_fe2o3_core::prelude::*;
11use crate::{
12 time::{CalClock, CalClockZone},
13 schedule::{RecurrencePattern, Action},
14};
15use std::{
16 fmt,
17 sync::atomic::{AtomicU64, Ordering},
18};
19
20#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
21pub struct TaskId(u64);
22
23impl TaskId {
24 fn next() -> Self {
25 static COUNTER: AtomicU64 = AtomicU64::new(1);
26 TaskId(COUNTER.fetch_add(1, Ordering::Relaxed))
27 }
28}
29
30impl 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)]
37pub 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)]
47pub enum TaskPriority {
48 Low = 1,
49 Normal = 2,
50 High = 3,
51 Critical = 4,
52}
53
54impl Default for TaskPriority {
55 fn default() -> Self {
56 TaskPriority::Normal
57 }
58}
59
60#[derive(Debug, Clone)]
61pub 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
69impl 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)]
82pub 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
99impl 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)]
121pub 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
130impl 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
198pub 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
209impl 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)]
316mod 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}