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 | |
| 10 | use oxedyne_fe2o3_core::prelude::*; |
| 11 | use crate::{ |
| 12 | schedule::{Task, TaskId, TaskExecutor}, |
| 13 | time::CalClock, |
| 14 | }; |
| 15 | use 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)] |
| 24 | pub 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 | |
| 33 | impl 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)] |
| 47 | pub 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)] |
| 58 | struct QueuedTask { |
| 59 | task: Task, |
| 60 | queued_at: Instant, |
| 61 | priority_score: u64, |
| 62 | } |
| 63 | |
| 64 | impl PartialEq for QueuedTask { |
| 65 | fn eq(&self, other: &Self) -> bool { |
| 66 | self.priority_score == other.priority_score |
| 67 | } |
| 68 | } |
| 69 | |
| 70 | impl Eq for QueuedTask {} |
| 71 | |
| 72 | impl PartialOrd for QueuedTask { |
| 73 | fn partial_cmp(&self, other: &Self) -> Option<Ordering> { |
| 74 | Some(self.cmp(other)) |
| 75 | } |
| 76 | } |
| 77 | |
| 78 | impl 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 | |
| 86 | impl 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)] |
| 137 | enum 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 | |
| 147 | pub 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 | |
| 160 | impl 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 | |
| 173 | impl 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 | |
| 529 | impl Drop for Scheduler { |
| 530 | fn drop(&mut self) { |
| 531 | let _ = self.stop(); |
| 532 | } |
| 533 | } |