use crate::metrics::{compute_report, ScoreConfig};
use crate::observation::Observation;
use crate::schema::StabilityReport;
pub struct StreamingScorer {
config: ScoreConfig,
current_id: Option<String>,
buffer: Vec<Observation>,
}
impl StreamingScorer {
pub fn new(config: ScoreConfig) -> Self {
Self {
config,
current_id: None,
buffer: Vec::new(),
}
}
pub fn push(&mut self, obs: Observation) -> Option<StabilityReport> {
if self.current_id.as_deref() == Some(obs.sample_id.as_str()) {
self.buffer.push(obs);
None
} else {
let result = self.flush_inner();
self.current_id = Some(obs.sample_id.clone());
self.buffer.push(obs);
result
}
}
pub fn flush(&mut self) -> Option<StabilityReport> {
self.flush_inner()
}
fn flush_inner(&mut self) -> Option<StabilityReport> {
if self.buffer.is_empty() {
return None;
}
let id = self.current_id.take().unwrap();
let report = compute_report(&id, &self.buffer, &self.config);
self.buffer.clear();
Some(report)
}
}