oxicode_sdk/ports/inmem/
resources.rs1use std::future::Future;
4use std::pin::Pin;
5use std::sync::Arc;
6use std::sync::atomic::{AtomicU64, Ordering};
7
8use crate::SdkError;
9use crate::ports::{ResourceMonitor, ResourceUsage};
10
11pub struct CountingResourceMonitor {
21 active_agents: Arc<AtomicU64>,
22 tokens: Arc<AtomicU64>,
23}
24
25impl std::fmt::Debug for CountingResourceMonitor {
26 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
27 f.debug_struct("CountingResourceMonitor").finish()
28 }
29}
30
31impl Default for CountingResourceMonitor {
32 fn default() -> Self {
33 Self::new()
34 }
35}
36
37impl CountingResourceMonitor {
38 pub fn new() -> Self {
40 Self {
41 active_agents: Arc::new(AtomicU64::new(0)),
42 tokens: Arc::new(AtomicU64::new(0)),
43 }
44 }
45
46 pub fn inc_active(&self) {
48 self.active_agents.fetch_add(1, Ordering::Relaxed);
49 }
50
51 pub fn dec_active(&self) {
53 self.active_agents.fetch_sub(1, Ordering::Relaxed);
54 }
55
56 pub fn add_tokens(&self, n: u64) {
58 self.tokens.fetch_add(n, Ordering::Relaxed);
59 }
60}
61
62impl ResourceMonitor for CountingResourceMonitor {
63 fn snapshot(
64 &self,
65 ) -> Pin<Box<dyn Future<Output = Result<ResourceUsage, SdkError>> + Send + '_>> {
66 let active = self.active_agents.load(Ordering::Relaxed) as usize;
67 let tokens = self.tokens.load(Ordering::Relaxed);
68 Box::pin(async move {
69 Ok(ResourceUsage {
70 cpu_percent: 0.0,
71 memory_bytes: 0,
72 disk_bytes: 0,
73 active_agents: active,
74 tokens_consumed: tokens,
75 })
76 })
77 }
78}
79
80#[cfg(test)]
81mod tests {
82 use super::*;
83
84 #[tokio::test]
85 async fn counters_increment() {
86 let m = CountingResourceMonitor::new();
87 m.inc_active();
88 m.inc_active();
89 m.add_tokens(100);
90 let snap = m.snapshot().await.unwrap();
91 assert_eq!(snap.active_agents, 2);
92 assert_eq!(snap.tokens_consumed, 100);
93 m.dec_active();
94 let snap = m.snapshot().await.unwrap();
95 assert_eq!(snap.active_agents, 1);
96 }
97}