rill_patchbay/
observer.rs1use parking_lot::RwLock;
4use rill_core::queues::telemetry::Telemetry;
5use rill_core::traits::ParameterId;
6use rill_core_actor::ActorRef;
7use std::collections::HashMap;
8use std::sync::Arc;
9use std::time::{SystemTime, UNIX_EPOCH};
10
11#[derive(Debug, Clone, Default)]
13pub struct ComponentStats {
14 pub operations: u64,
16 pub total_time_ns: u64,
18 pub max_time_ns: u64,
20 pub violations: u64,
22 pub avg_time_ns: f64,
24}
25
26#[derive(Debug, Clone)]
28pub struct Violation {
29 pub component: String,
31 pub expected_ns: u64,
33 pub actual_ns: u64,
35 pub timestamp: u64,
37 pub value: Option<f32>,
39}
40
41#[derive(Debug, Default, Clone)]
43pub struct SandboxSummary {
44 pub total_operations: u64,
46 pub total_violations: u64,
48 pub components: Vec<String>,
50 pub max_time_ns: u64,
52 pub max_time_component: Option<String>,
54 pub violations_count: usize,
56}
57
58#[derive(Clone)]
60pub struct MicroControlObserver {
61 stats: Arc<RwLock<HashMap<String, ComponentStats>>>,
62 violations: Arc<RwLock<Vec<Violation>>>,
63 telemetry_tx: Option<ActorRef<Telemetry>>,
64}
65
66impl Default for MicroControlObserver {
67 fn default() -> Self {
68 Self::new()
69 }
70}
71
72impl MicroControlObserver {
73 pub fn new() -> Self {
75 Self {
76 stats: Arc::new(RwLock::new(HashMap::new())),
77 violations: Arc::new(RwLock::new(Vec::new())),
78 telemetry_tx: None,
79 }
80 }
81
82 pub fn with_actor(tx: ActorRef<Telemetry>) -> Self {
84 Self {
85 stats: Arc::new(RwLock::new(HashMap::new())),
86 violations: Arc::new(RwLock::new(Vec::new())),
87 telemetry_tx: Some(tx),
88 }
89 }
90
91 fn send_telemetry(&self, event: Telemetry) {
92 if let Some(ref tx) = self.telemetry_tx {
93 tx.send(event);
94 }
95 }
96
97 pub fn observe_start(&self, component: &str) -> OperationGuard {
99 OperationGuard {
100 component: component.to_string(),
101 start_time: Self::now(),
102 observer: self.clone(),
103 }
104 }
105
106 pub fn observe_start_with_params(
108 &self,
109 component: &str,
110 _port: String,
111 _parameter: &ParameterId,
112 ) -> OperationGuard {
113 let guard = self.observe_start(component);
114 self.send_telemetry(Telemetry::event(
115 "observer",
116 "micro_start",
117 vec![0.0_f32, 0.0_f32],
118 ));
119 guard
120 }
121
122 pub fn record_violation(
124 &self,
125 component: &str,
126 expected_ns: u64,
127 actual_ns: u64,
128 value: Option<f32>,
129 ) {
130 let violation = Violation {
131 component: component.to_string(),
132 expected_ns,
133 actual_ns,
134 timestamp: Self::now(),
135 value,
136 };
137 self.violations.write().push(violation.clone());
138 let mut stats = self.stats.write();
139 let comp_stats = stats.entry(component.to_string()).or_default();
140 comp_stats.violations += 1;
141 self.send_telemetry(Telemetry::violation(
142 component,
143 expected_ns,
144 actual_ns,
145 value,
146 ));
147 }
148
149 pub fn component_stats(&self, component: &str) -> Option<ComponentStats> {
151 self.stats.read().get(component).cloned()
152 }
153
154 pub fn violations(&self) -> Vec<Violation> {
156 self.violations.read().clone()
157 }
158
159 pub fn sandbox_summary(&self) -> SandboxSummary {
161 let stats = self.stats.read();
162 let mut summary = SandboxSummary::default();
163 for (component, comp_stats) in stats.iter() {
164 summary.total_operations += comp_stats.operations;
165 summary.total_violations += comp_stats.violations;
166 summary.components.push(component.clone());
167 if comp_stats.max_time_ns > summary.max_time_ns {
168 summary.max_time_ns = comp_stats.max_time_ns;
169 summary.max_time_component = Some(component.clone());
170 }
171 }
172 summary.violations_count = self.violations.read().len();
173 summary
174 }
175
176 fn now() -> u64 {
177 SystemTime::now()
178 .duration_since(UNIX_EPOCH)
179 .unwrap_or_default()
180 .as_micros() as u64
181 }
182}
183
184pub struct OperationGuard {
186 component: String,
187 start_time: u64,
188 observer: MicroControlObserver,
189}
190
191impl Drop for OperationGuard {
192 fn drop(&mut self) {
193 let duration = (Self::now() - self.start_time) * 1000;
194 let mut stats = self.observer.stats.write();
195 let comp_stats = stats.entry(self.component.clone()).or_default();
196 comp_stats.operations += 1;
197 comp_stats.total_time_ns += duration;
198 if duration > comp_stats.max_time_ns {
199 comp_stats.max_time_ns = duration;
200 }
201 comp_stats.avg_time_ns = comp_stats.total_time_ns as f64 / comp_stats.operations as f64;
202 self.observer.send_telemetry(Telemetry::event(
203 "observer",
204 "micro_complete",
205 vec![duration as f32],
206 ));
207 }
208}
209
210impl OperationGuard {
211 fn now() -> u64 {
212 SystemTime::now()
213 .duration_since(UNIX_EPOCH)
214 .unwrap_or_default()
215 .as_micros() as u64
216 }
217}
218
219#[cfg(test)]
220mod tests {
221 use super::*;
222 use rill_core_actor::ActorSystem;
223 use std::sync::{Arc, Mutex};
224
225 #[test]
226 fn test_observer_creation() {
227 let observer = MicroControlObserver::new();
228 let stats = observer.sandbox_summary();
229 assert_eq!(stats.total_operations, 0);
230 }
231
232 #[test]
233 fn test_observer_record_violation() {
234 let system = ActorSystem::new();
235 let received = Arc::new(Mutex::new(Vec::new()));
236 let recv = received.clone();
237 let mut actor = system.spawn("telemetry", move |msg: Telemetry| {
238 recv.lock().unwrap().push(msg);
239 });
240 let observer = MicroControlObserver::with_actor(actor.actor_ref());
241 observer.record_violation("test_comp", 100, 250, Some(0.5));
242 let stats = observer.sandbox_summary();
243 assert_eq!(stats.total_violations, 1);
244 let violations = observer.violations();
245 assert_eq!(violations.len(), 1);
246 assert_eq!(violations[0].component, "test_comp");
247 actor.drain();
248 let events = received.lock().unwrap();
249 for evt in events.iter() {
250 if let Telemetry::Violation { component, .. } = evt {
251 assert_eq!(component, "test_comp");
252 }
253 }
254 }
255
256 #[test]
257 fn test_observer_operation_guard() {
258 let system = ActorSystem::new();
259 let received = Arc::new(Mutex::new(Vec::new()));
260 let recv = received.clone();
261 let mut actor = system.spawn("telemetry", move |msg: Telemetry| {
262 recv.lock().unwrap().push(msg);
263 });
264 let observer = MicroControlObserver::with_actor(actor.actor_ref());
265 {
266 let _guard = observer.observe_start("test_op");
267 std::thread::sleep(std::time::Duration::from_micros(10));
268 }
269 let stats = observer.sandbox_summary();
270 assert_eq!(stats.total_operations, 1);
271 actor.drain();
272 let events = received.lock().unwrap();
273 for evt in events.iter() {
274 if let Telemetry::Event { kind, .. } = evt {
275 assert_eq!(kind, "micro_complete");
276 }
277 }
278 }
279}