Skip to main content

corescout_substrate/
reflector.rs

1//! The reflector: the loop that produces `M(t)` from a real machine.
2//!
3//! # Why this is called a reflector and not a mirror
4//!
5//! `corescout-mirror` holds the *reflection*: a snapshot, and the shared memory
6//! it is published through. Both are inert data. This is the machinery that
7//! produces one, and it is the part that needs hardware, sensors and
8//! privileges.
9//!
10//! Keeping the two apart is what lets a consumer link the reflection without
11//! linking the thing that makes it. An observer holding a `MirrorSnapshot` has
12//! no `Reflector`, no sensors and no `/sys`.
13//!
14//! # One pass
15//!
16//! ```text
17//! clear the matrix        no stale value survives into a new reflection
18//! stamp the clocks
19//! for each bound sensor:
20//!     time it, run it, record what it cost and how confident it is
21//! ```
22//!
23//! Clearing first is the anti-staleness rule. A sensor that fails this tick
24//! leaves a hole with a reason attached, not last minute's temperature wearing
25//! a fresh timestamp. A consumer cannot tell those apart, so the mirror must
26//! not offer it the chance to be wrong.
27
28use 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
38/// A live producer of reflections.
39///
40/// One `Reflector` owns one machine's observation loop. Calling
41/// [`Reflector::observe`] advances it to the present; nothing else changes it.
42/// No consumer, model or agent can write to a reflection. The only way to
43/// change what the mirror shows is to change the machine.
44pub 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    /// Discover structure, then bind every sensor to it.
74    ///
75    /// Binding is where all the once-per-epoch work happens: paths are
76    /// constructed, capabilities probed, descriptors opened, entities and
77    /// channels declared. A sensor that cannot bind is not an error. It is
78    /// recorded as inactive **with a reason**, and the mirror says so rather
79    /// than implying this machine has no thermal sensors.
80    pub fn build(substrate: Substrate, sensors: Vec<Box<dyn Sensor>>) -> Result<Reflector> {
81        Reflector::build_at_epoch(substrate, sensors, 1)
82    }
83
84    /// Build at a given epoch, for a rebuild after the machine's shape changed.
85    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                    // Why it could not bind is the useful part. "Permission
121                    // denied" is actionable; "unsupported" never will be.
122                    let availability = classify(&error, descriptor.requires_privilege);
123                    // A sensor that failed to bind must not leave half-declared
124                    // columns behind.
125                    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    /// Advance the reflection to the present.
176    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        // Columns belonging to a sensor that could not bind are marked once,
184        // with that sensor's reason, so a consumer sees "unsupported on this
185        // machine" rather than an undifferentiated hole.
186        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            // The cost of looking is itself observed, and published.
209            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    /// The current reflection, as an owned snapshot.
229    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    /// Total time the last observation pass spent looking.
246    pub fn last_observation_cost_ns(&self) -> u64 {
247        self.reports.iter().map(|r| r.last_cost_ns).sum()
248    }
249
250    /// Sensors that could not bind on this machine, and why.
251    pub fn inactive_sensors(&self) -> Vec<&SensorReport> {
252        self.reports.iter().filter(|r| r.inactive).collect()
253    }
254
255    /// Whether the machine's shape has changed since this reflector was built.
256    ///
257    /// Cheap enough to call every tick. When it returns true the caller should
258    /// rebuild at a higher epoch, because every row index it holds is stale.
259    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            // If we cannot tell, do not claim the machine changed: a spurious
270            // epoch bump discards every consumer's cached indices.
271            Err(_) => false,
272        }
273    }
274}
275
276/// Work out why a sensor could not bind.
277fn 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            // A sensor that needs privilege and reports "unsupported" is
287            // ambiguous on a machine where the interface exists but is
288            // unreadable. Attribute it to privilege, which is the actionable
289            // reading, and let the sensor say otherwise if it knows better.
290            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
307/// Wall-clock nanoseconds since the Unix epoch.
308///
309/// For correlating a reflection with logs and with other machines. Never for
310/// measuring intervals: it is not monotonic, and `monotonic_ns` exists for that.
311fn 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    /// A sensor that observes nothing real. Enough to exercise the loop.
329    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        // A cell no sensor claims is "not applicable", not "unknown".
414        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        // The distinction that tells an operator whether running as root would
439        // help.
440        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}