Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_datime/src/schedule/scheduler.rs

17.8 KiB, 41 runs

created by r1870400018:8496, 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//! The scheduler itself.
2//!
3//! Tasks sit in a priority queue, highest score first and FIFO within a score.
4//! A background thread moves them to worker threads as their time arrives, or
5//! the caller can drive the whole thing by hand with process_pending.
6//!
7//! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\
8//! Anthropic Claude
9
10use oxedyne_fe2o3_core::prelude::*;
11use crate::{
12 schedule::{Task, TaskId, TaskExecutor},
13 time::CalClock,
14};
15use std::{
16 collections::{HashMap, BinaryHeap},
17 cmp::Ordering,
18 sync::{Arc, Mutex, mpsc},
19 thread::{self, JoinHandle},
20 time::{Duration, Instant},
21};
22
23#[derive(Debug, Clone)]
24pub struct SchedulerConfig {
25 pub max_concurrent_tasks: usize,
26 pub check_interval_millis: u64,
27 pub continue_on_failure: bool,
28 pub queue_size: usize, // zero means unlimited
29 pub worker_threads: usize,
30 pub enable_background_processing: bool,
31}
32
33impl Default for SchedulerConfig {
34 fn default() -> Self {
35 SchedulerConfig {
36 max_concurrent_tasks: 10,
37 check_interval_millis: 1000, // 1 second
38 continue_on_failure: true,
39 queue_size: 1000,
40 worker_threads: 4,
41 enable_background_processing: true,
42 }
43 }
44}
45
46#[derive(Debug, Clone, Default)]
47pub struct SchedulerStats {
48 pub scheduled_tasks: usize,
49 pub running_tasks: usize,
50 pub queued_tasks: usize,
51 pub completed_tasks: u64,
52 pub failed_tasks: u64,
53 pub avg_execution_time_millis: u64,
54 pub uptime_seconds: u64,
55}
56
57#[derive(Debug)]
58struct QueuedTask {
59 task: Task,
60 queued_at: Instant,
61 priority_score: u64,
62}
63
64impl PartialEq for QueuedTask {
65 fn eq(&self, other: &Self) -> bool {
66 self.priority_score == other.priority_score
67 }
68}
69
70impl Eq for QueuedTask {}
71
72impl PartialOrd for QueuedTask {
73 fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
74 Some(self.cmp(other))
75 }
76}
77
78impl Ord for QueuedTask {
79 fn cmp(&self, other: &Self) -> Ordering {
80 // Higher priority scores come first (reverse order for max-heap behaviour)
81 other.priority_score.cmp(&self.priority_score)
82 .then_with(|| self.queued_at.cmp(&other.queued_at)) // FIFO for same priority
83 }
84}
85
86impl QueuedTask {
87 fn new(task: Task) -> Self {
88 let priority_score = Self::calculate_priority_score(&task);
89 Self {
90 task,
91 queued_at: Instant::now(),
92 priority_score,
93 }
94 }
95
96 fn calculate_priority_score(task: &Task) -> u64 {
97 let base_priority = match task.priority {
98 crate::schedule::TaskPriority::Critical => 1000,
99 crate::schedule::TaskPriority::High => 750,
100 crate::schedule::TaskPriority::Normal => 500,
101 crate::schedule::TaskPriority::Low => 250,
102 };
103
104 // Adjust for deadline urgency - use a safe fallback if clock operations fail
105 let now = CalClock::now_utc()
106 .or_else(|_| CalClock::new(2024, 1, 1, 0, 0, 0, 0, crate::time::CalClockZone::utc()))
107 .unwrap_or_else(|_| {
108 // Ultimate fallback - create a minimal clock
109 task.scheduled_time.clone()
110 });
111
112 let time_until_scheduled = if task.scheduled_time >= now {
113 let task_millis = task.scheduled_time.to_nanos_since_epoch().unwrap_or(0) / 1_000_000;
114 let now_millis = now.to_nanos_since_epoch().unwrap_or(0) / 1_000_000;
115 (task_millis - now_millis).max(0) as u64
116 } else {
117 // Overdue tasks get maximum urgency
118 return base_priority + 10000;
119 };
120
121 // Closer deadlines get higher priority
122 let urgency_bonus = if time_until_scheduled < 60000 { // < 1 minute
123 500
124 } else if time_until_scheduled < 300000 { // < 5 minutes
125 300
126 } else if time_until_scheduled < 3600000 { // < 1 hour
127 100
128 } else {
129 0
130 };
131
132 base_priority + urgency_bonus
133 }
134}
135
136#[derive(Debug)]
137enum SchedulerMessage {
138 Stop,
139 #[allow(dead_code)]
140 AddTask(Task),
141 #[allow(dead_code)]
142 RemoveTask(TaskId),
143 #[allow(dead_code)]
144 GetStats,
145}
146
147pub struct Scheduler {
148 config: SchedulerConfig,
149 tasks: HashMap<TaskId, Task>,
150 task_queue: Arc<Mutex<BinaryHeap<QueuedTask>>>,
151 stats: Arc<Mutex<SchedulerStats>>,
152 executor: Arc<Mutex<TaskExecutor>>,
153 worker_handles: Vec<JoinHandle<()>>,
154 message_sender: Option<mpsc::Sender<SchedulerMessage>>,
155 scheduler_handle: Option<JoinHandle<()>>,
156 start_time: Instant,
157 is_running: Arc<Mutex<bool>>,
158}
159
160impl std::fmt::Debug for Scheduler {
161 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
162 f.debug_struct("Scheduler")
163 .field("config", &self.config)
164 .field("task_count", &self.tasks.len())
165 .field("is_running", &self.is_running.lock().map(|guard| *guard).unwrap_or_else(|_| {
166 // For Debug impl, provide a reasonable fallback for poisoned mutex
167 false
168 }))
169 .finish()
170 }
171}
172
173impl Scheduler {
174 pub fn new() -> Self {
175 Self::with_config(SchedulerConfig::default())
176 }
177
178 pub fn with_config(config: SchedulerConfig) -> Self {
179 Scheduler {
180 config,
181 tasks: HashMap::new(),
182 task_queue: Arc::new(Mutex::new(BinaryHeap::new())),
183 stats: Arc::new(Mutex::new(SchedulerStats::default())),
184 executor: Arc::new(Mutex::new(TaskExecutor::new())),
185 worker_handles: Vec::new(),
186 message_sender: None,
187 scheduler_handle: None,
188 start_time: Instant::now(),
189 is_running: Arc::new(Mutex::new(false)),
190 }
191 }
192
193 pub fn start(&mut self) -> Outcome<()> {
194 if *lock_mutex!(self.is_running) {
195 return Err(err!("Scheduler is already running"; Invalid, Input));
196 }
197
198 *lock_mutex!(self.is_running) = true;
199 self.start_time = Instant::now();
200
201 if self.config.enable_background_processing {
202 // Start the main scheduler thread
203 let (tx, rx) = mpsc::channel();
204 self.message_sender = Some(tx);
205
206 let queue = Arc::clone(&self.task_queue);
207 let stats = Arc::clone(&self.stats);
208 let executor = Arc::clone(&self.executor);
209 let is_running = Arc::clone(&self.is_running);
210 let config = self.config.clone();
211
212 self.scheduler_handle = Some(thread::spawn(move || {
213 Self::scheduler_loop(rx, queue, stats, executor, is_running, config);
214 }));
215
216 // Start worker threads
217 for worker_id in 0..self.config.worker_threads {
218 let queue = Arc::clone(&self.task_queue);
219 let stats = Arc::clone(&self.stats);
220 let executor = Arc::clone(&self.executor);
221 let is_running = Arc::clone(&self.is_running);
222
223 let handle = thread::spawn(move || {
224 Self::worker_loop(worker_id, queue, stats, executor, is_running);
225 });
226
227 self.worker_handles.push(handle);
228 }
229 }
230
231 Ok(())
232 }
233
234 pub fn stop(&mut self) -> Outcome<()> {
235 *lock_mutex!(self.is_running) = false;
236
237 // Send stop message to scheduler thread
238 if let Some(sender) = &self.message_sender {
239 let _ = sender.send(SchedulerMessage::Stop);
240 }
241
242 // Wait for scheduler thread to finish
243 if let Some(handle) = self.scheduler_handle.take() {
244 let _ = handle.join();
245 }
246
247 // Wait for all worker threads to finish
248 for handle in self.worker_handles.drain(..) {
249 let _ = handle.join();
250 }
251
252 self.message_sender = None;
253 Ok(())
254 }
255
256 pub fn schedule(&mut self, task: Task) -> Outcome<TaskId> {
257 let task_id = task.id;
258
259 // Add to local task registry
260 self.tasks.insert(task_id, task.clone());
261
262 // Add to processing queue if background processing is enabled
263 if self.config.enable_background_processing && self.message_sender.is_some() {
264 let queued_task = QueuedTask::new(task);
265 {
266 let mut queue = lock_mutex!(self.task_queue);
267 if self.config.queue_size == 0 || queue.len() < self.config.queue_size {
268 queue.push(queued_task);
269 } else {
270 return Err(err!("Task queue is full"; Invalid, Input));
271 }
272 }
273 }
274
275 // Update statistics
276 {
277 let mut stats = lock_mutex!(self.stats);
278 stats.scheduled_tasks = self.tasks.len();
279 stats.queued_tasks = lock_mutex!(self.task_queue).len();
280 }
281
282 Ok(task_id)
283 }
284
285 pub fn unschedule(&mut self, task_id: TaskId) -> Outcome<()> {
286 self.tasks.remove(&task_id);
287
288 // Update statistics
289 {
290 let mut stats = lock_mutex!(self.stats);
291 stats.scheduled_tasks = self.tasks.len();
292 }
293
294 Ok(())
295 }
296
297 /// The manual alternative to background processing, for callers that
298 /// would rather drive the scheduler themselves.
299 pub fn process_pending(&mut self) -> Outcome<()> {
300 if self.config.enable_background_processing {
301 return Err(err!("Cannot manually process when background processing is enabled"; Invalid, Input));
302 }
303
304 let now = res!(CalClock::now_utc());
305 let mut tasks_to_execute = Vec::new();
306
307 // Find tasks ready for execution
308 for (task_id, task) in &self.tasks {
309 if task.is_ready_to_execute(&now) {
310 tasks_to_execute.push(*task_id);
311 }
312 }
313
314 // Execute ready tasks
315 let mut stats = lock_mutex!(self.stats);
316 let mut executor = lock_mutex!(self.executor);
317
318 for task_id in tasks_to_execute {
319 if let Some(mut task) = self.tasks.remove(&task_id) {
320 match executor.execute_with_retry(&mut task) {
321 Ok(result) => {
322 if result.success {
323 stats.completed_tasks += 1;
324 } else {
325 stats.failed_tasks += 1;
326 }
327
328 // Handle recurring tasks
329 if task.recurrence.is_some() {
330 if let Ok(()) = task.advance_to_next_execution() {
331 self.tasks.insert(task_id, task);
332 }
333 }
334 },
335 Err(e) => {
336 eprintln!("Task execution error: {}", e);
337 stats.failed_tasks += 1;
338 }
339 }
340 }
341 }
342
343 stats.scheduled_tasks = self.tasks.len();
344 Ok(())
345 }
346
347 fn scheduler_loop(
348 rx: mpsc::Receiver<SchedulerMessage>,
349 queue: Arc<Mutex<BinaryHeap<QueuedTask>>>,
350 stats: Arc<Mutex<SchedulerStats>>,
351 _executor: Arc<Mutex<TaskExecutor>>,
352 is_running: Arc<Mutex<bool>>,
353 config: SchedulerConfig,
354 ) {
355 let check_interval = Duration::from_millis(config.check_interval_millis);
356
357 while {
358 match is_running.lock() {
359 Ok(guard) => *guard,
360 Err(_) => {
361 eprintln!("Scheduler thread: poisoned is_running mutex, stopping");
362 false
363 }
364 }
365 } {
366 // Process any incoming messages
367 while let Ok(message) = rx.try_recv() {
368 match message {
369 SchedulerMessage::Stop => {
370 *lock_mutex_thread!(is_running, "setting is_running to false") = false;
371 return;
372 },
373 SchedulerMessage::AddTask(task) => {
374 let queued_task = QueuedTask::new(task);
375 let mut queue = lock_mutex_thread!(queue, "scheduler thread queue access");
376 queue.push(queued_task);
377 },
378 SchedulerMessage::RemoveTask(_task_id) => {
379 // TODO: Implement task removal from queue
380 },
381 SchedulerMessage::GetStats => {
382 // TODO: Implement stats reporting
383 }
384 }
385 }
386
387 // Update statistics
388 {
389 let mut stats = lock_mutex_thread!(stats, "stats update");
390 stats.queued_tasks = lock_mutex_thread!(queue, "queue length check").len();
391 }
392
393 thread::sleep(check_interval);
394 }
395 }
396
397 fn worker_loop(
398 _worker_id: usize,
399 queue: Arc<Mutex<BinaryHeap<QueuedTask>>>,
400 stats: Arc<Mutex<SchedulerStats>>,
401 executor: Arc<Mutex<TaskExecutor>>,
402 is_running: Arc<Mutex<bool>>,
403 ) {
404 while {
405 match is_running.lock() {
406 Ok(guard) => *guard,
407 Err(_) => {
408 eprintln!("Scheduler thread: poisoned is_running mutex, stopping");
409 false
410 }
411 }
412 } {
413 // Try to get a task from the queue
414 let task_opt = {
415 let mut queue = lock_mutex_thread!(queue, "queue access");
416 queue.pop()
417 };
418
419 if let Some(mut queued_task) = task_opt {
420 let now = CalClock::now_utc()
421 .or_else(|_| CalClock::new(2024, 1, 1, 0, 0, 0, 0, crate::time::CalClockZone::utc()))
422 .unwrap_or_else(|_| queued_task.task.scheduled_time.clone());
423
424 // Check if task is ready to execute
425 if queued_task.task.is_ready_to_execute(&now) {
426 // Update running task count
427 {
428 let mut stats = lock_mutex_thread!(stats, "stats update");
429 stats.running_tasks += 1;
430 }
431
432 // Execute the task
433 let execution_result = {
434 let mut executor = lock_mutex_thread!(executor, "executor access");
435 executor.execute_with_retry(&mut queued_task.task)
436 };
437
438 // Update statistics
439 {
440 let mut stats = lock_mutex_thread!(stats, "stats update");
441 stats.running_tasks = stats.running_tasks.saturating_sub(1);
442
443 match execution_result {
444 Ok(result) => {
445 if result.success {
446 stats.completed_tasks += 1;
447 } else {
448 stats.failed_tasks += 1;
449 }
450 // Update average execution time
451 let total_tasks = stats.completed_tasks + stats.failed_tasks;
452 if total_tasks > 0 {
453 stats.avg_execution_time_millis =
454 (stats.avg_execution_time_millis * (total_tasks - 1) + result.duration_millis) / total_tasks;
455 }
456 },
457 Err(_) => {
458 stats.failed_tasks += 1;
459 }
460 }
461 }
462
463 // Handle recurring tasks
464 if queued_task.task.recurrence.is_some() {
465 if let Ok(()) = queued_task.task.advance_to_next_execution() {
466 // Re-queue the task for next execution
467 let mut queue = lock_mutex_thread!(queue, "scheduler thread queue access");
468 queue.push(queued_task);
469 }
470 }
471 } else {
472 // Task not ready yet, put it back in the queue
473 let mut queue = lock_mutex_thread!(queue, "queue access");
474 queue.push(queued_task);
475 drop(queue);
476
477 // Sleep briefly to avoid busy waiting
478 thread::sleep(Duration::from_millis(100));
479 }
480 } else {
481 // No tasks available, sleep briefly
482 thread::sleep(Duration::from_millis(100));
483 }
484 }
485 }
486
487 pub fn stats(&self) -> SchedulerStats {
488 let mut stats = match self.stats.lock() {
489 Ok(guard) => guard.clone(),
490 Err(_) => {
491 eprintln!("Warning: Stats mutex poisoned, returning default stats");
492 return SchedulerStats::default();
493 }
494 };
495 stats.uptime_seconds = self.start_time.elapsed().as_secs();
496 stats.scheduled_tasks = self.tasks.len();
497 stats
498 }
499
500 pub fn get_task(&self, task_id: TaskId) -> Option<&Task> {
501 self.tasks.get(&task_id)
502 }
503
504 pub fn list_tasks(&self) -> Vec<&Task> {
505 self.tasks.values().collect()
506 }
507
508 pub fn queue_size(&self) -> usize {
509 match self.task_queue.lock() {
510 Ok(guard) => guard.len(),
511 Err(_) => {
512 eprintln!("Warning: Task queue mutex poisoned, returning 0");
513 0
514 }
515 }
516 }
517
518 pub fn is_running(&self) -> bool {
519 match self.is_running.lock() {
520 Ok(guard) => *guard,
521 Err(_) => {
522 eprintln!("Warning: Is_running mutex poisoned, returning false");
523 false
524 }
525 }
526 }
527}
528
529impl Drop for Scheduler {
530 fn drop(&mut self) {
531 let _ = self.stop();
532 }
533}