1use corescout_core::clock;
29use corescout_core::Result;
30use corescout_mirror::schema::{Availability, AvailabilityMatrix, SensorReport};
31use corescout_mirror::{
32 ChannelSpec, Entity, MirrorSnapshot, Relation, StateMatrix, FORMAT_VERSION,
33};
34
35use crate::discovery::Substrate;
36use crate::observation::{AvailabilityWriter, BindContext, Sensor, StateWriter};
37
38pub struct Reflector {
45 substrate: Substrate,
46 sensors: Vec<Box<dyn Sensor>>,
47 entities: Vec<Entity>,
48 channels: Vec<ChannelSpec>,
49 relations: Vec<Relation>,
50 state: StateMatrix,
51 availability: AvailabilityMatrix,
52 reports: Vec<SensorReport>,
53 epoch: u64,
54 sequence: u64,
55 monotonic_ns: u64,
56 realtime_ns: u64,
57}
58
59impl std::fmt::Debug for Reflector {
60 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
61 f.debug_struct("Reflector")
62 .field("entities", &self.entities.len())
63 .field("channels", &self.channels.len())
64 .field("relations", &self.relations.len())
65 .field("sensors", &self.sensors.len())
66 .field("epoch", &self.epoch)
67 .field("sequence", &self.sequence)
68 .finish()
69 }
70}
71
72impl Reflector {
73 pub fn build(substrate: Substrate, sensors: Vec<Box<dyn Sensor>>) -> Result<Reflector> {
81 Reflector::build_at_epoch(substrate, sensors, 1)
82 }
83
84 pub fn build_at_epoch(
86 substrate: Substrate,
87 sensors: Vec<Box<dyn Sensor>>,
88 epoch: u64,
89 ) -> Result<Reflector> {
90 let (mut entities, mut relations) = substrate.structure();
91 let mut channels: Vec<ChannelSpec> = Vec::new();
92 let mut reports = Vec::new();
93 let mut bound = Vec::new();
94
95 for mut sensor in sensors {
96 let descriptor = sensor.descriptor();
97 debug_assert!(
98 descriptor.perturbation.is_passive(),
99 "sensor `{}` declares Material perturbation and belongs in experiment/",
100 descriptor.key
101 );
102
103 let before = channels.len();
104 let outcome = {
105 let mut ctx = BindContext::new(
106 &substrate,
107 descriptor.id,
108 &mut entities,
109 &mut channels,
110 &mut relations,
111 );
112 sensor.bind(&mut ctx)
113 };
114
115 reports.push(match outcome {
116 Ok(()) => {
117 SensorReport::pending(descriptor.id, descriptor.key, descriptor.perturbation)
118 }
119 Err(error) => {
120 let availability = classify(&error, descriptor.requires_privilege);
123 channels.truncate(before);
126 SensorReport::unavailable(
127 descriptor.id,
128 descriptor.key,
129 descriptor.perturbation,
130 availability,
131 )
132 }
133 });
134 bound.push(sensor);
135 }
136
137 let state = StateMatrix::new(entities.len(), channels.len());
138 let availability = AvailabilityMatrix::new(entities.len(), channels.len());
139 Ok(Reflector {
140 substrate,
141 sensors: bound,
142 entities,
143 channels,
144 relations,
145 state,
146 availability,
147 reports,
148 epoch,
149 sequence: 0,
150 monotonic_ns: 0,
151 realtime_ns: 0,
152 })
153 }
154
155 pub fn substrate(&self) -> &Substrate {
156 &self.substrate
157 }
158
159 pub fn epoch(&self) -> u64 {
160 self.epoch
161 }
162
163 pub fn entities(&self) -> &[Entity] {
164 &self.entities
165 }
166
167 pub fn channels(&self) -> &[ChannelSpec] {
168 &self.channels
169 }
170
171 pub fn relations(&self) -> &[Relation] {
172 &self.relations
173 }
174
175 pub fn observe(&mut self) {
177 self.state.clear();
178 self.availability.clear();
179 self.monotonic_ns = clock::now_ns();
180 self.realtime_ns = realtime_ns();
181 self.sequence += 1;
182
183 for report in self.reports.iter().filter(|r| r.inactive) {
187 for (col, channel) in self.channels.iter().enumerate() {
188 if channel.sensor == report.id {
189 for row in 0..self.availability.rows() {
190 self.availability.set(row, col, report.availability);
191 }
192 }
193 }
194 }
195
196 for (sensor, report) in self.sensors.iter_mut().zip(self.reports.iter_mut()) {
197 if report.inactive {
198 continue;
199 }
200 let started = clock::now_ns();
201 let outcome = {
202 let mut writer = StateWriter::new(
203 &mut self.state,
204 AvailabilityWriter::new(&mut self.availability),
205 );
206 sensor.observe(&mut writer)
207 };
208 report.last_cost_ns = clock::now_ns().saturating_sub(started);
210 report.samples = outcome.samples;
211 report.errors = outcome.errors;
212 report.sample_age_ns = outcome.sample_age_ns;
213 report.sampling_latency_ns = report.last_cost_ns;
214 let attempted = outcome.samples + outcome.errors;
215 report.confidence = if attempted == 0 {
216 0.0
217 } else {
218 outcome.samples as f64 / attempted as f64
219 };
220 report.availability = if outcome.samples > 0 {
221 Availability::Observed
222 } else {
223 Availability::Unavailable
224 };
225 }
226 }
227
228 pub fn snapshot(&self) -> MirrorSnapshot {
230 MirrorSnapshot {
231 format_version: FORMAT_VERSION,
232 epoch: self.epoch,
233 sequence: self.sequence,
234 monotonic_ns: self.monotonic_ns,
235 realtime_ns: self.realtime_ns,
236 entities: self.entities.clone(),
237 channels: self.channels.clone(),
238 relations: self.relations.clone(),
239 state: self.state.clone(),
240 availability: self.availability.clone(),
241 sensors: self.reports.clone(),
242 }
243 }
244
245 pub fn last_observation_cost_ns(&self) -> u64 {
247 self.reports.iter().map(|r| r.last_cost_ns).sum()
248 }
249
250 pub fn inactive_sensors(&self) -> Vec<&SensorReport> {
252 self.reports.iter().filter(|r| r.inactive).collect()
253 }
254
255 pub fn shape_changed(&self) -> bool {
260 match self.substrate.rediscover() {
261 Ok(current) => {
262 let (entities, _) = current.structure();
263 entities.len() != self.entities.len()
264 || entities
265 .iter()
266 .zip(&self.entities)
267 .any(|(now, before)| now.id != before.id)
268 }
269 Err(_) => false,
272 }
273 }
274}
275
276fn classify(error: &corescout_core::Error, requires_privilege: bool) -> Availability {
278 use corescout_core::Error;
279 match error {
280 Error::Io { source, .. } => match source.kind() {
281 std::io::ErrorKind::PermissionDenied => Availability::PermissionDenied,
282 std::io::ErrorKind::NotFound => Availability::Unsupported,
283 _ => Availability::Unavailable,
284 },
285 Error::Unsupported(_) => {
286 if requires_privilege {
291 Availability::PermissionDenied
292 } else {
293 Availability::Unsupported
294 }
295 }
296 Error::Syscall { errno, .. } => {
297 if *errno == 1 || *errno == 13 {
298 Availability::PermissionDenied
299 } else {
300 Availability::Unavailable
301 }
302 }
303 _ => Availability::Unknown,
304 }
305}
306
307fn realtime_ns() -> u64 {
312 std::time::SystemTime::now()
313 .duration_since(std::time::UNIX_EPOCH)
314 .map(|d| d.as_nanos() as u64)
315 .unwrap_or(0)
316}
317
318#[cfg(test)]
319mod tests {
320 use super::*;
321 use crate::discovery::Roots;
322 use crate::observation::{Perturbation, SensorDescriptor, SensorOutcome, Uncertainty};
323 use crate::test_support::fake_topology;
324 use corescout_mirror::schema::SensorId;
325 use corescout_mirror::state::{Semantics, Unit};
326 use corescout_mirror::ChannelId;
327
328 struct FakeSensor {
330 rows: Vec<u32>,
331 channel: Option<ChannelId>,
332 bind_error: Option<corescout_core::Error>,
333 }
334
335 impl FakeSensor {
336 fn working() -> FakeSensor {
337 FakeSensor {
338 rows: Vec::new(),
339 channel: None,
340 bind_error: None,
341 }
342 }
343
344 fn failing(error: corescout_core::Error) -> FakeSensor {
345 FakeSensor {
346 rows: Vec::new(),
347 channel: None,
348 bind_error: Some(error),
349 }
350 }
351 }
352
353 impl Sensor for FakeSensor {
354 fn descriptor(&self) -> SensorDescriptor {
355 SensorDescriptor {
356 id: SensorId(900),
357 key: "fake",
358 physical_fact: "nothing; this sensor exists for tests",
359 source: "none",
360 max_rate_hz: 1000.0,
361 perturbation: Perturbation::None,
362 uncertainty: Uncertainty::unknown("not a real measurement"),
363 requires_privilege: false,
364 }
365 }
366
367 fn bind(&mut self, ctx: &mut BindContext<'_>) -> Result<()> {
368 if let Some(error) = self.bind_error.take() {
369 return Err(error);
370 }
371 self.channel =
372 Some(ctx.declare_channel("fake.value", Unit::Dimensionless, Semantics::Instant));
373 for cpu in &ctx.substrate().topology.logical_cpus {
374 if let Some(row) = ctx.row_of(&corescout_mirror::entity::keys::logical_cpu(cpu.id))
375 {
376 self.rows.push(row);
377 }
378 }
379 Ok(())
380 }
381
382 fn observe(&mut self, out: &mut StateWriter<'_>) -> SensorOutcome {
383 let mut outcome = SensorOutcome::default();
384 let Some(channel) = self.channel else {
385 return outcome;
386 };
387 for row in &self.rows {
388 out.set(*row, channel, 1.0);
389 outcome.sample();
390 }
391 outcome
392 }
393 }
394
395 fn substrate() -> Substrate {
396 Substrate::new(fake_topology(), Roots::new("/nonexistent", "/nonexistent"))
397 }
398
399 #[test]
400 fn observing_fills_the_matrix_and_marks_it_observed() {
401 let mut reflector =
402 Reflector::build(substrate(), vec![Box::new(FakeSensor::working())]).unwrap();
403 reflector.observe();
404 let snapshot = reflector.snapshot();
405
406 assert_eq!(snapshot.sequence, 1);
407 assert_eq!(snapshot.lookup("cpu/0", "fake.value"), Some(1.0));
408
409 let row = snapshot.row_of_key("cpu/0").unwrap();
410 let col = snapshot.channel("fake.value").unwrap();
411 assert_eq!(snapshot.availability(row, col), Availability::Observed);
412
413 let machine = snapshot.row_of_key("machine").unwrap();
415 assert_eq!(
416 snapshot.availability(machine, col),
417 Availability::NotApplicable
418 );
419 }
420
421 #[test]
422 fn a_sensor_that_cannot_bind_reports_why() {
423 let reflector = Reflector::build(
424 substrate(),
425 vec![Box::new(FakeSensor::failing(
426 corescout_core::Error::unsupported("no such interface"),
427 ))],
428 )
429 .unwrap();
430 let inactive = reflector.inactive_sensors();
431 assert_eq!(inactive.len(), 1);
432 assert_eq!(inactive[0].availability, Availability::Unsupported);
433 assert_eq!(inactive[0].confidence, 0.0);
434 }
435
436 #[test]
437 fn a_permission_failure_is_distinguished_from_an_absent_interface() {
438 let denied = corescout_core::Error::io(
441 "/sys/class/powercap/intel-rapl:0/energy_uj",
442 std::io::Error::from(std::io::ErrorKind::PermissionDenied),
443 );
444 let reflector =
445 Reflector::build(substrate(), vec![Box::new(FakeSensor::failing(denied))]).unwrap();
446 assert_eq!(
447 reflector.inactive_sensors()[0].availability,
448 Availability::PermissionDenied
449 );
450 }
451
452 #[test]
453 fn each_pass_starts_from_a_clean_matrix() {
454 let mut reflector =
455 Reflector::build(substrate(), vec![Box::new(FakeSensor::working())]).unwrap();
456 reflector.observe();
457 assert!(reflector.snapshot().state.observed_cells() > 0);
458
459 reflector.sensors.clear();
460 reflector.observe();
461 assert_eq!(
462 reflector.snapshot().state.observed_cells(),
463 0,
464 "stale values must not survive into a later reflection"
465 );
466 }
467
468 #[test]
469 fn the_reflector_records_what_looking_cost() {
470 let mut reflector =
471 Reflector::build(substrate(), vec![Box::new(FakeSensor::working())]).unwrap();
472 reflector.observe();
473 let snapshot = reflector.snapshot();
474 assert_eq!(snapshot.sensors.len(), 1);
475 assert!(snapshot.sensors[0].samples > 0);
476 assert_eq!(snapshot.sensors[0].confidence, 1.0);
477 assert_eq!(
478 snapshot.observation_cost_ns(),
479 reflector.last_observation_cost_ns()
480 );
481 }
482}