use corescout_core::clock;
use corescout_core::Result;
use corescout_mirror::schema::{Availability, AvailabilityMatrix, SensorReport};
use corescout_mirror::{
ChannelSpec, Entity, MirrorSnapshot, Relation, StateMatrix, FORMAT_VERSION,
};
use crate::discovery::Substrate;
use crate::observation::{AvailabilityWriter, BindContext, Sensor, StateWriter};
pub struct Reflector {
substrate: Substrate,
sensors: Vec<Box<dyn Sensor>>,
entities: Vec<Entity>,
channels: Vec<ChannelSpec>,
relations: Vec<Relation>,
state: StateMatrix,
availability: AvailabilityMatrix,
reports: Vec<SensorReport>,
epoch: u64,
sequence: u64,
monotonic_ns: u64,
realtime_ns: u64,
}
impl std::fmt::Debug for Reflector {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Reflector")
.field("entities", &self.entities.len())
.field("channels", &self.channels.len())
.field("relations", &self.relations.len())
.field("sensors", &self.sensors.len())
.field("epoch", &self.epoch)
.field("sequence", &self.sequence)
.finish()
}
}
impl Reflector {
pub fn build(substrate: Substrate, sensors: Vec<Box<dyn Sensor>>) -> Result<Reflector> {
Reflector::build_at_epoch(substrate, sensors, 1)
}
pub fn build_at_epoch(
substrate: Substrate,
sensors: Vec<Box<dyn Sensor>>,
epoch: u64,
) -> Result<Reflector> {
let (mut entities, mut relations) = substrate.structure();
let mut channels: Vec<ChannelSpec> = Vec::new();
let mut reports = Vec::new();
let mut bound = Vec::new();
for mut sensor in sensors {
let descriptor = sensor.descriptor();
debug_assert!(
descriptor.perturbation.is_passive(),
"sensor `{}` declares Material perturbation and belongs in experiment/",
descriptor.key
);
let before = channels.len();
let outcome = {
let mut ctx = BindContext::new(
&substrate,
descriptor.id,
&mut entities,
&mut channels,
&mut relations,
);
sensor.bind(&mut ctx)
};
reports.push(match outcome {
Ok(()) => {
SensorReport::pending(descriptor.id, descriptor.key, descriptor.perturbation)
}
Err(error) => {
let availability = classify(&error, descriptor.requires_privilege);
channels.truncate(before);
SensorReport::unavailable(
descriptor.id,
descriptor.key,
descriptor.perturbation,
availability,
)
}
});
bound.push(sensor);
}
let state = StateMatrix::new(entities.len(), channels.len());
let availability = AvailabilityMatrix::new(entities.len(), channels.len());
Ok(Reflector {
substrate,
sensors: bound,
entities,
channels,
relations,
state,
availability,
reports,
epoch,
sequence: 0,
monotonic_ns: 0,
realtime_ns: 0,
})
}
pub fn substrate(&self) -> &Substrate {
&self.substrate
}
pub fn epoch(&self) -> u64 {
self.epoch
}
pub fn entities(&self) -> &[Entity] {
&self.entities
}
pub fn channels(&self) -> &[ChannelSpec] {
&self.channels
}
pub fn relations(&self) -> &[Relation] {
&self.relations
}
pub fn observe(&mut self) {
self.state.clear();
self.availability.clear();
self.monotonic_ns = clock::now_ns();
self.realtime_ns = realtime_ns();
self.sequence += 1;
for report in self.reports.iter().filter(|r| r.inactive) {
for (col, channel) in self.channels.iter().enumerate() {
if channel.sensor == report.id {
for row in 0..self.availability.rows() {
self.availability.set(row, col, report.availability);
}
}
}
}
for (sensor, report) in self.sensors.iter_mut().zip(self.reports.iter_mut()) {
if report.inactive {
continue;
}
let started = clock::now_ns();
let outcome = {
let mut writer = StateWriter::new(
&mut self.state,
AvailabilityWriter::new(&mut self.availability),
);
sensor.observe(&mut writer)
};
report.last_cost_ns = clock::now_ns().saturating_sub(started);
report.samples = outcome.samples;
report.errors = outcome.errors;
report.sample_age_ns = outcome.sample_age_ns;
report.sampling_latency_ns = report.last_cost_ns;
let attempted = outcome.samples + outcome.errors;
report.confidence = if attempted == 0 {
0.0
} else {
outcome.samples as f64 / attempted as f64
};
report.availability = if outcome.samples > 0 {
Availability::Observed
} else {
Availability::Unavailable
};
}
}
pub fn snapshot(&self) -> MirrorSnapshot {
MirrorSnapshot {
format_version: FORMAT_VERSION,
epoch: self.epoch,
sequence: self.sequence,
monotonic_ns: self.monotonic_ns,
realtime_ns: self.realtime_ns,
entities: self.entities.clone(),
channels: self.channels.clone(),
relations: self.relations.clone(),
state: self.state.clone(),
availability: self.availability.clone(),
sensors: self.reports.clone(),
}
}
pub fn last_observation_cost_ns(&self) -> u64 {
self.reports.iter().map(|r| r.last_cost_ns).sum()
}
pub fn inactive_sensors(&self) -> Vec<&SensorReport> {
self.reports.iter().filter(|r| r.inactive).collect()
}
pub fn shape_changed(&self) -> bool {
match self.substrate.rediscover() {
Ok(current) => {
let (entities, _) = current.structure();
entities.len() != self.entities.len()
|| entities
.iter()
.zip(&self.entities)
.any(|(now, before)| now.id != before.id)
}
Err(_) => false,
}
}
}
fn classify(error: &corescout_core::Error, requires_privilege: bool) -> Availability {
use corescout_core::Error;
match error {
Error::Io { source, .. } => match source.kind() {
std::io::ErrorKind::PermissionDenied => Availability::PermissionDenied,
std::io::ErrorKind::NotFound => Availability::Unsupported,
_ => Availability::Unavailable,
},
Error::Unsupported(_) => {
if requires_privilege {
Availability::PermissionDenied
} else {
Availability::Unsupported
}
}
Error::Syscall { errno, .. } => {
if *errno == 1 || *errno == 13 {
Availability::PermissionDenied
} else {
Availability::Unavailable
}
}
_ => Availability::Unknown,
}
}
fn realtime_ns() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_nanos() as u64)
.unwrap_or(0)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::discovery::Roots;
use crate::observation::{Perturbation, SensorDescriptor, SensorOutcome, Uncertainty};
use crate::test_support::fake_topology;
use corescout_mirror::schema::SensorId;
use corescout_mirror::state::{Semantics, Unit};
use corescout_mirror::ChannelId;
struct FakeSensor {
rows: Vec<u32>,
channel: Option<ChannelId>,
bind_error: Option<corescout_core::Error>,
}
impl FakeSensor {
fn working() -> FakeSensor {
FakeSensor {
rows: Vec::new(),
channel: None,
bind_error: None,
}
}
fn failing(error: corescout_core::Error) -> FakeSensor {
FakeSensor {
rows: Vec::new(),
channel: None,
bind_error: Some(error),
}
}
}
impl Sensor for FakeSensor {
fn descriptor(&self) -> SensorDescriptor {
SensorDescriptor {
id: SensorId(900),
key: "fake",
physical_fact: "nothing; this sensor exists for tests",
source: "none",
max_rate_hz: 1000.0,
perturbation: Perturbation::None,
uncertainty: Uncertainty::unknown("not a real measurement"),
requires_privilege: false,
}
}
fn bind(&mut self, ctx: &mut BindContext<'_>) -> Result<()> {
if let Some(error) = self.bind_error.take() {
return Err(error);
}
self.channel =
Some(ctx.declare_channel("fake.value", Unit::Dimensionless, Semantics::Instant));
for cpu in &ctx.substrate().topology.logical_cpus {
if let Some(row) = ctx.row_of(&corescout_mirror::entity::keys::logical_cpu(cpu.id))
{
self.rows.push(row);
}
}
Ok(())
}
fn observe(&mut self, out: &mut StateWriter<'_>) -> SensorOutcome {
let mut outcome = SensorOutcome::default();
let Some(channel) = self.channel else {
return outcome;
};
for row in &self.rows {
out.set(*row, channel, 1.0);
outcome.sample();
}
outcome
}
}
fn substrate() -> Substrate {
Substrate::new(fake_topology(), Roots::new("/nonexistent", "/nonexistent"))
}
#[test]
fn observing_fills_the_matrix_and_marks_it_observed() {
let mut reflector =
Reflector::build(substrate(), vec![Box::new(FakeSensor::working())]).unwrap();
reflector.observe();
let snapshot = reflector.snapshot();
assert_eq!(snapshot.sequence, 1);
assert_eq!(snapshot.lookup("cpu/0", "fake.value"), Some(1.0));
let row = snapshot.row_of_key("cpu/0").unwrap();
let col = snapshot.channel("fake.value").unwrap();
assert_eq!(snapshot.availability(row, col), Availability::Observed);
let machine = snapshot.row_of_key("machine").unwrap();
assert_eq!(
snapshot.availability(machine, col),
Availability::NotApplicable
);
}
#[test]
fn a_sensor_that_cannot_bind_reports_why() {
let reflector = Reflector::build(
substrate(),
vec![Box::new(FakeSensor::failing(
corescout_core::Error::unsupported("no such interface"),
))],
)
.unwrap();
let inactive = reflector.inactive_sensors();
assert_eq!(inactive.len(), 1);
assert_eq!(inactive[0].availability, Availability::Unsupported);
assert_eq!(inactive[0].confidence, 0.0);
}
#[test]
fn a_permission_failure_is_distinguished_from_an_absent_interface() {
let denied = corescout_core::Error::io(
"/sys/class/powercap/intel-rapl:0/energy_uj",
std::io::Error::from(std::io::ErrorKind::PermissionDenied),
);
let reflector =
Reflector::build(substrate(), vec![Box::new(FakeSensor::failing(denied))]).unwrap();
assert_eq!(
reflector.inactive_sensors()[0].availability,
Availability::PermissionDenied
);
}
#[test]
fn each_pass_starts_from_a_clean_matrix() {
let mut reflector =
Reflector::build(substrate(), vec![Box::new(FakeSensor::working())]).unwrap();
reflector.observe();
assert!(reflector.snapshot().state.observed_cells() > 0);
reflector.sensors.clear();
reflector.observe();
assert_eq!(
reflector.snapshot().state.observed_cells(),
0,
"stale values must not survive into a later reflection"
);
}
#[test]
fn the_reflector_records_what_looking_cost() {
let mut reflector =
Reflector::build(substrate(), vec![Box::new(FakeSensor::working())]).unwrap();
reflector.observe();
let snapshot = reflector.snapshot();
assert_eq!(snapshot.sensors.len(), 1);
assert!(snapshot.sensors[0].samples > 0);
assert_eq!(snapshot.sensors[0].confidence, 1.0);
assert_eq!(
snapshot.observation_cost_ns(),
reflector.last_observation_cost_ns()
);
}
}