1use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
11use parking_lot::RwLock;
12
13use super::{NodeResources, Workload, ResourceRequirements};
14use crate::types::NodeId;
15
16#[derive(Debug)]
18struct NodeScoreCache {
19 node_id: NodeId,
21 cpu_score: u32,
23 memory_score: u32,
25 gpu_score: u32,
27 combined_score: u32,
29 cpu_available: u64,
31 memory_available: u64,
33 gpu_available: u32,
35 schedulable: bool,
37}
38
39impl NodeScoreCache {
40 fn from_node(node: &NodeResources) -> Self {
41 let cpu_available = node.cpu_available();
42 let memory_available = node.memory_available();
43 let gpu_available = node.gpus_available() as u32;
44
45 let cpu_score = ((cpu_available as f64 / node.cpu_capacity.max(1) as f64) * 1000.0) as u32;
47 let memory_score = ((memory_available as f64 / node.memory_capacity.max(1) as f64) * 1000.0) as u32;
48 let gpu_score = if node.gpus.is_empty() {
49 500
50 } else {
51 ((gpu_available as f64 / node.gpus.len() as f64) * 1000.0) as u32
52 };
53
54 let combined_score = (cpu_score + memory_score + gpu_score) / 3;
56
57 Self {
58 node_id: node.node_id,
59 cpu_score,
60 memory_score,
61 gpu_score,
62 combined_score,
63 cpu_available,
64 memory_available,
65 gpu_available,
66 schedulable: node.schedulable,
67 }
68 }
69
70 #[inline(always)]
71 fn can_fit(&self, req: &ResourceRequirements) -> bool {
72 self.schedulable
73 && self.cpu_available >= req.cpu_millis
74 && self.memory_available >= req.memory_mb
75 && self.gpu_available >= req.gpu_count
76 }
77
78 #[inline(always)]
79 fn score_for_workload(&self, req: &ResourceRequirements) -> u32 {
80 if !self.can_fit(req) {
81 return 0;
82 }
83
84 let cpu_fit = 1000 - ((self.cpu_available - req.cpu_millis) * 1000 / self.cpu_available.max(1)) as u32;
87 let mem_fit = 1000 - ((self.memory_available - req.memory_mb) * 1000 / self.memory_available.max(1)) as u32;
88
89 (cpu_fit * 4 + mem_fit * 4 + self.gpu_score * 2) / 10
91 }
92
93 #[inline]
101 fn commit(&mut self, req: &ResourceRequirements, cpu_capacity: u64, memory_capacity: u64, gpu_total: u32) {
102 self.cpu_available = self.cpu_available.saturating_sub(req.cpu_millis);
103 self.memory_available = self.memory_available.saturating_sub(req.memory_mb);
104 self.gpu_available = self.gpu_available.saturating_sub(req.gpu_count);
105 self.recompute_scores(cpu_capacity, memory_capacity, gpu_total);
106 }
107
108 #[inline]
112 fn release(&mut self, req: &ResourceRequirements, cpu_capacity: u64, memory_capacity: u64, gpu_total: u32) {
113 self.cpu_available = (self.cpu_available + req.cpu_millis).min(cpu_capacity);
114 self.memory_available = (self.memory_available + req.memory_mb).min(memory_capacity);
115 self.gpu_available = (self.gpu_available + req.gpu_count).min(gpu_total);
116 self.recompute_scores(cpu_capacity, memory_capacity, gpu_total);
117 }
118
119 #[inline]
122 fn recompute_scores(&mut self, cpu_capacity: u64, memory_capacity: u64, gpu_total: u32) {
123 self.cpu_score = ((self.cpu_available as f64 / cpu_capacity.max(1) as f64) * 1000.0) as u32;
124 self.memory_score = ((self.memory_available as f64 / memory_capacity.max(1) as f64) * 1000.0) as u32;
125 self.gpu_score = if gpu_total == 0 {
126 500
127 } else {
128 ((self.gpu_available as f64 / gpu_total as f64) * 1000.0) as u32
129 };
130 self.combined_score = (self.cpu_score + self.memory_score + self.gpu_score) / 3;
131 }
132}
133
134pub struct WorkloadBatch {
136 workloads: Vec<Workload>,
137 results: Vec<Option<NodeId>>,
138}
139
140impl WorkloadBatch {
141 pub fn new(workloads: Vec<Workload>) -> Self {
143 let len = workloads.len();
144 Self {
145 workloads,
146 results: vec![None; len],
147 }
148 }
149
150 pub fn results(&self) -> &[Option<NodeId>] {
152 &self.results
153 }
154
155 pub fn workloads(&self) -> &[Workload] {
157 &self.workloads
158 }
159}
160
161pub struct OptimizedScheduler {
169 node_cache: RwLock<Vec<NodeScoreCache>>,
171 nodes: RwLock<Vec<NodeResources>>,
173 scheduled_count: AtomicU64,
175 total_time_ns: AtomicU64,
177 cache_generation: AtomicUsize,
179}
180
181impl OptimizedScheduler {
182 pub fn new() -> Self {
184 Self {
185 node_cache: RwLock::new(Vec::new()),
186 nodes: RwLock::new(Vec::new()),
187 scheduled_count: AtomicU64::new(0),
188 total_time_ns: AtomicU64::new(0),
189 cache_generation: AtomicUsize::new(0),
190 }
191 }
192
193 pub fn register_node(&self, node: NodeResources) {
195 let cache = NodeScoreCache::from_node(&node);
196 self.nodes.write().push(node);
197 self.node_cache.write().push(cache);
198 self.cache_generation.fetch_add(1, Ordering::Relaxed);
199 }
200
201 pub fn refresh_cache(&self) {
203 let nodes = self.nodes.read();
204 let mut cache = self.node_cache.write();
205 cache.clear();
206 cache.extend(nodes.iter().map(NodeScoreCache::from_node));
207 self.cache_generation.fetch_add(1, Ordering::Relaxed);
208 }
209
210 #[inline]
221 pub fn schedule_fast(&self, workload: &Workload) -> Option<NodeId> {
222 let start = std::time::Instant::now();
223 let cache = self.node_cache.read();
224
225 if cache.is_empty() {
226 return None;
227 }
228
229 let req = &workload.resources;
230
231 let best = cache
238 .iter()
239 .filter(|n| n.can_fit(req))
240 .max_by_key(|n| n.score_for_workload(req))
241 .map(|n| n.node_id);
242
243 self.scheduled_count.fetch_add(1, Ordering::Relaxed);
245 self.total_time_ns.fetch_add(start.elapsed().as_nanos() as u64, Ordering::Relaxed);
246
247 best
248 }
249
250 pub fn schedule_fast_commit(&self, workload: &Workload) -> Option<NodeId> {
271 let start = std::time::Instant::now();
272 let req = &workload.resources;
273
274 let mut nodes = self.nodes.write();
276 let mut cache = self.node_cache.write();
277
278 if cache.is_empty() {
279 return None;
280 }
281
282 let best_idx = {
285 let scored = cache.iter().enumerate().filter(|(_, n)| n.can_fit(req));
286 scored.max_by_key(|(_, n)| n.score_for_workload(req)).map(|(i, _)| i)
290 };
291
292 let result = match best_idx {
293 Some(idx) => {
294 let node = &mut nodes[idx];
298 if node.allocate(req) {
299 let node_id = node.node_id;
300 let cpu_capacity = node.cpu_capacity;
301 let memory_capacity = node.memory_capacity;
302 let gpu_total = node.gpus.len() as u32;
303
304 cache[idx].commit(req, cpu_capacity, memory_capacity, gpu_total);
306
307 Some(node_id)
308 } else {
309 None
310 }
311 }
312 None => None,
313 };
314
315 drop(cache);
318 drop(nodes);
319
320 self.scheduled_count.fetch_add(1, Ordering::Relaxed);
321 self.total_time_ns.fetch_add(start.elapsed().as_nanos() as u64, Ordering::Relaxed);
322
323 result
324 }
325
326 pub fn release_workload(&self, node_id: NodeId, req: &ResourceRequirements, gpu_ids: &[u32]) {
338 let mut nodes = self.nodes.write();
339 let mut cache = self.node_cache.write();
340
341 if let Some(idx) = nodes.iter().position(|n| n.node_id == node_id) {
342 let node = &mut nodes[idx];
343 node.release(req, gpu_ids);
344 let cpu_capacity = node.cpu_capacity;
345 let memory_capacity = node.memory_capacity;
346 let gpu_total = node.gpus.len() as u32;
347 if let Some(entry) = cache.get_mut(idx) {
348 entry.release(req, cpu_capacity, memory_capacity, gpu_total);
349 }
350 }
351 }
352
353 pub fn schedule_batch(&self, batch: &mut WorkloadBatch) {
361 let start = std::time::Instant::now();
362 let cache = self.node_cache.read();
363
364 if cache.is_empty() {
365 return;
366 }
367
368 let mut indices: Vec<usize> = (0..batch.workloads.len()).collect();
370 indices.sort_by(|&a, &b| {
371 batch.workloads[b].priority.cmp(&batch.workloads[a].priority)
372 });
373
374 let mut node_allocated: Vec<(u64, u64, u32)> = cache.iter()
376 .map(|n| (n.cpu_available, n.memory_available, n.gpu_available))
377 .collect();
378
379 for idx in indices {
381 let workload = &batch.workloads[idx];
382 let req = &workload.resources;
383
384 let mut best_node: Option<usize> = None;
386 let mut best_score: u32 = 0;
387
388 for (i, (n, alloc)) in cache.iter().zip(node_allocated.iter()).enumerate() {
389 if !n.schedulable {
390 continue;
391 }
392
393 if alloc.0 < req.cpu_millis || alloc.1 < req.memory_mb || alloc.2 < req.gpu_count {
395 continue;
396 }
397
398 let remaining_cpu = alloc.0 - req.cpu_millis;
400 let remaining_mem = alloc.1 - req.memory_mb;
401
402 let score = 2000 - (remaining_cpu * 1000 / n.cpu_available.max(1)) as u32
404 - (remaining_mem * 1000 / n.memory_available.max(1)) as u32;
405
406 if score > best_score {
407 best_score = score;
408 best_node = Some(i);
409 }
410 }
411
412 if let Some(node_idx) = best_node {
413 batch.results[idx] = Some(cache[node_idx].node_id);
414
415 node_allocated[node_idx].0 -= req.cpu_millis;
417 node_allocated[node_idx].1 -= req.memory_mb;
418 node_allocated[node_idx].2 -= req.gpu_count;
419 }
420 }
421
422 let count = batch.workloads.len() as u64;
424 self.scheduled_count.fetch_add(count, Ordering::Relaxed);
425 self.total_time_ns.fetch_add(start.elapsed().as_nanos() as u64, Ordering::Relaxed);
426 }
427
428 pub fn commit_batch(&self, batch: &mut WorkloadBatch) {
439 let mut nodes = self.nodes.write();
440 let mut cache = self.node_cache.write();
441
442 if nodes.is_empty() {
443 return;
444 }
445
446 for i in 0..batch.workloads.len() {
448 let Some(node_id) = batch.results[i] else { continue };
449 let Some(idx) = nodes.iter().position(|n| n.node_id == node_id) else {
450 batch.results[i] = None;
451 continue;
452 };
453
454 let req = &batch.workloads[i].resources;
455 let node = &mut nodes[idx];
456 if node.allocate(req) {
457 let cpu_capacity = node.cpu_capacity;
458 let memory_capacity = node.memory_capacity;
459 let gpu_total = node.gpus.len() as u32;
460 if let Some(entry) = cache.get_mut(idx) {
461 entry.commit(req, cpu_capacity, memory_capacity, gpu_total);
462 }
463 } else {
464 batch.results[i] = None;
466 }
467 }
468 }
469
470 pub fn schedule_and_commit_batch(&self, batch: &mut WorkloadBatch) {
476 self.schedule_batch(batch);
477 self.commit_batch(batch);
478 }
479
480 pub fn stats(&self) -> SchedulerStats {
482 let count = self.scheduled_count.load(Ordering::Relaxed);
483 let time_ns = self.total_time_ns.load(Ordering::Relaxed);
484
485 SchedulerStats {
486 total_scheduled: count,
487 total_time_ns: time_ns,
488 avg_time_ns: if count > 0 { time_ns / count } else { 0 },
489 decisions_per_sec: if time_ns > 0 {
490 (count as f64 * 1_000_000_000.0 / time_ns as f64) as u64
491 } else {
492 0
493 },
494 node_count: self.node_cache.read().len(),
495 }
496 }
497
498 pub fn reset_stats(&self) {
500 self.scheduled_count.store(0, Ordering::Relaxed);
501 self.total_time_ns.store(0, Ordering::Relaxed);
502 }
503
504 pub fn node_count(&self) -> usize {
506 self.node_cache.read().len()
507 }
508
509 pub fn utilization(&self) -> ClusterUtilization {
511 let nodes = self.nodes.read();
512
513 let mut total_cpu: u64 = 0;
514 let mut used_cpu: u64 = 0;
515 let mut total_mem: u64 = 0;
516 let mut used_mem: u64 = 0;
517 let mut total_gpu: u32 = 0;
518 let mut used_gpu: u32 = 0;
519
520 for node in nodes.iter() {
521 total_cpu += node.cpu_capacity;
522 used_cpu += node.cpu_allocated;
523 total_mem += node.memory_capacity;
524 used_mem += node.memory_allocated;
525 total_gpu += node.gpus.len() as u32;
526 used_gpu += node.gpus_allocated.len() as u32;
527 }
528
529 ClusterUtilization {
530 cpu_percent: if total_cpu > 0 { (used_cpu as f64 / total_cpu as f64) * 100.0 } else { 0.0 },
531 memory_percent: if total_mem > 0 { (used_mem as f64 / total_mem as f64) * 100.0 } else { 0.0 },
532 gpu_percent: if total_gpu > 0 { (used_gpu as f64 / total_gpu as f64) * 100.0 } else { 0.0 },
533 total_cpu,
534 used_cpu,
535 total_memory: total_mem,
536 used_memory: used_mem,
537 total_gpus: total_gpu,
538 used_gpus: used_gpu,
539 }
540 }
541}
542
543impl Default for OptimizedScheduler {
544 fn default() -> Self {
545 Self::new()
546 }
547}
548
549#[derive(Debug, Clone)]
551pub struct SchedulerStats {
552 pub total_scheduled: u64,
554 pub total_time_ns: u64,
556 pub avg_time_ns: u64,
558 pub decisions_per_sec: u64,
560 pub node_count: usize,
562}
563
564#[derive(Debug, Clone)]
566pub struct ClusterUtilization {
567 pub cpu_percent: f64,
569 pub memory_percent: f64,
571 pub gpu_percent: f64,
573 pub total_cpu: u64,
575 pub used_cpu: u64,
577 pub total_memory: u64,
579 pub used_memory: u64,
581 pub total_gpus: u32,
583 pub used_gpus: u32,
585}
586
587pub struct FFDBinPacker {
591 nodes: Vec<NodeResources>,
593}
594
595impl FFDBinPacker {
596 pub fn new(mut nodes: Vec<NodeResources>) -> Self {
598 nodes.sort_by(|a, b| {
600 let cap_a = a.cpu_capacity + a.memory_capacity;
601 let cap_b = b.cpu_capacity + b.memory_capacity;
602 cap_b.cmp(&cap_a)
603 });
604 Self { nodes }
605 }
606
607 pub fn pack(&mut self, mut workloads: Vec<Workload>) -> (Vec<(String, NodeId)>, f64) {
610 workloads.sort_by(|a, b| {
612 let req_a = a.resources.cpu_millis + a.resources.memory_mb;
613 let req_b = b.resources.cpu_millis + b.resources.memory_mb;
614 req_b.cmp(&req_a)
615 });
616
617 let mut assignments = Vec::new();
618 let mut node_usage: Vec<(u64, u64)> = self.nodes.iter()
619 .map(|n| (0u64, 0u64))
620 .collect();
621
622 for workload in &workloads {
623 let req = &workload.resources;
624
625 for (i, node) in self.nodes.iter().enumerate() {
627 let (used_cpu, used_mem) = node_usage[i];
628 let avail_cpu = node.cpu_capacity.saturating_sub(used_cpu);
629 let avail_mem = node.memory_capacity.saturating_sub(used_mem);
630
631 if avail_cpu >= req.cpu_millis && avail_mem >= req.memory_mb {
632 assignments.push((workload.id.clone(), node.node_id));
633 node_usage[i].0 += req.cpu_millis;
634 node_usage[i].1 += req.memory_mb;
635 break;
636 }
637 }
638 }
639
640 let total_cpu: u64 = self.nodes.iter().map(|n| n.cpu_capacity).sum();
642 let used_cpu: u64 = node_usage.iter().map(|(c, _)| c).sum();
643 let utilization = if total_cpu > 0 {
644 (used_cpu as f64 / total_cpu as f64) * 100.0
645 } else {
646 0.0
647 };
648
649 (assignments, utilization)
650 }
651}
652
653#[cfg(test)]
654mod tests {
655 use super::*;
656
657 fn create_nodes(count: usize) -> Vec<NodeResources> {
658 (0..count).map(|_| {
659 NodeResources::new(NodeId::new(), 8000, 32768)
660 }).collect()
661 }
662
663 fn create_workloads(count: usize) -> Vec<Workload> {
664 (0..count).map(|i| {
665 Workload::new(format!("w-{}", i), "test")
666 .with_resources(ResourceRequirements::new()
667 .cpu(100 + (i as u64 % 10) * 100)
668 .memory(256 + (i as u64 % 8) * 256))
669 }).collect()
670 }
671
672 #[test]
673 fn test_optimized_scheduler_fast() {
674 let scheduler = OptimizedScheduler::new();
675
676 for node in create_nodes(100) {
677 scheduler.register_node(node);
678 }
679
680 let workloads = create_workloads(1000);
681 let mut scheduled = 0;
682
683 for workload in &workloads {
684 if scheduler.schedule_fast(workload).is_some() {
685 scheduled += 1;
686 }
687 }
688
689 assert_eq!(scheduled, workloads.len(), "every workload should find a placement");
694
695 let stats = scheduler.stats();
696 println!("Scheduled: {}, Rate: {} decisions/sec", scheduled, stats.decisions_per_sec);
697 assert_eq!(stats.total_scheduled, workloads.len() as u64);
698 }
699
700 #[test]
701 fn test_batch_scheduling() {
702 let scheduler = OptimizedScheduler::new();
703
704 for node in create_nodes(50) {
705 scheduler.register_node(node);
706 }
707
708 let workloads = create_workloads(100);
709 let mut batch = WorkloadBatch::new(workloads);
710
711 scheduler.schedule_batch(&mut batch);
712
713 let scheduled: usize = batch.results().iter().filter(|r| r.is_some()).count();
714 assert!(scheduled > 0);
715 println!("Batch scheduled: {}/100", scheduled);
716 }
717
718 #[test]
719 fn test_schedule_fast_commit_reduces_capacity() {
720 let scheduler = OptimizedScheduler::new();
722 scheduler.register_node(NodeResources::new(NodeId::new(), 4000, 8192));
723
724 let req = ResourceRequirements::new().cpu(1000).memory(2048);
725 let workload = Workload::new("w", "test").with_resources(req.clone());
726
727 let mut placed = 0;
729 for _ in 0..6 {
730 if scheduler.schedule_fast_commit(&workload).is_some() {
731 placed += 1;
732 }
733 }
734 assert_eq!(placed, 4, "node should fit exactly 4 workloads after commit");
735
736 let util = scheduler.utilization();
738 assert_eq!(util.used_cpu, 4000);
739 assert_eq!(util.used_memory, 8192);
740
741 assert!(scheduler.schedule_fast(&workload).is_none());
744 }
745
746 #[test]
747 fn test_release_workload_restores_capacity() {
748 let scheduler = OptimizedScheduler::new();
749 let node_id = NodeId::new();
750 scheduler.register_node(NodeResources::new(node_id, 4000, 8192));
751
752 let req = ResourceRequirements::new().cpu(4000).memory(8192);
753 let workload = Workload::new("w", "test").with_resources(req.clone());
754
755 assert!(scheduler.schedule_fast_commit(&workload).is_some());
757 assert!(scheduler.schedule_fast_commit(&workload).is_none());
758
759 scheduler.release_workload(node_id, &req, &[]);
761 assert_eq!(scheduler.utilization().used_cpu, 0);
762 assert!(scheduler.schedule_fast_commit(&workload).is_some());
763 }
764
765 #[test]
766 fn test_commit_batch_writes_back() {
767 let scheduler = OptimizedScheduler::new();
768 for node in create_nodes(10) {
769 scheduler.register_node(node);
770 }
771
772 let workloads = create_workloads(50);
773 let mut batch = WorkloadBatch::new(workloads);
774
775 scheduler.schedule_and_commit_batch(&mut batch);
776
777 let scheduled: usize = batch.results().iter().filter(|r| r.is_some()).count();
778 assert!(scheduled > 0);
779
780 let util = scheduler.utilization();
782 assert!(util.used_cpu > 0, "commit_batch should write allocations to nodes");
783 }
784
785 #[test]
786 fn test_ffd_bin_packing() {
787 let nodes = create_nodes(10);
788 let workloads = create_workloads(50);
789
790 let mut packer = FFDBinPacker::new(nodes);
791 let (assignments, utilization) = packer.pack(workloads);
792
793 println!("FFD packed {} workloads, utilization: {:.1}%", assignments.len(), utilization);
794 assert!(assignments.len() > 0);
795 }
796}