///|
/// Streaming recording state for incremental RR ingestion.
pub(all) struct StreamingSession {
mut stats : OnlineStats
mut intervals : Array[Double]
config : HrvConfig
mut started : Bool
mut finished : Bool
} derive(Debug)
///|
/// A finished streaming session snapshot.
pub(all) struct StreamingSnapshot {
metrics : HrvMetrics
validation : IntervalValidation
diagnostic : SignalDiagnostic
sample_count : Int
duration_seconds : Double
finished : Bool
} derive(FromJson, ToJson, Debug, Eq)
///|
/// Create an empty streaming session.
pub fn StreamingSession::new(config : HrvConfig) -> StreamingSession {
{
stats: OnlineStats::new(),
intervals: [],
config,
started: false,
finished: false,
}
}
///|
/// Ingest one RR interval and update the online moments.
pub fn StreamingSession::push(
self : StreamingSession,
interval : Double,
) -> Unit {
if !self.finished {
self.started = true
self.intervals.push(interval)
self.stats.push(interval)
}
}
///|
/// Ingest a sequence of RR intervals.
pub fn StreamingSession::push_all(
self : StreamingSession,
values : Array[Double],
) -> Unit {
for value in values {
self.push(value)
}
}
///|
/// Mark a session finished; later pushes are ignored.
pub fn StreamingSession::finish(self : StreamingSession) -> Unit {
self.finished = true
}
///|
/// Return a snapshot without changing the session.
pub fn StreamingSession::snapshot(self : StreamingSession) -> StreamingSnapshot {
let validation = validate_intervals(self.intervals, self.config)
let quality = inspect_cleaning(self.intervals, Remove, self.config).quality
let metrics = calculate_metrics(self.intervals, self.config, quality)
{
metrics,
validation,
diagnostic: diagnose_signal(self.intervals, self.config),
sample_count: self.intervals.length(),
duration_seconds: sum_values(valid_intervals(self.intervals, self.config)) /
1000.0,
finished: self.finished,
}
}
///|
/// Return the current sample count.
pub fn StreamingSession::length(self : StreamingSession) -> Int {
self.intervals.length()
}
///|
/// Return whether the session is ready for recovery scoring.
pub fn StreamingSession::is_usable(self : StreamingSession) -> Bool {
is_recovery_signal_usable(self.snapshot().diagnostic)
}
///|
/// Return the latest interval, or zero before the first sample.
pub fn StreamingSession::last(self : StreamingSession) -> Double {
self.stats.last
}
///|
/// Return the current online mean without materializing a report.
pub fn StreamingSession::mean(self : StreamingSession) -> Double {
self.stats.mean
}
///|
/// Reset a session while retaining its configuration.
pub fn StreamingSession::reset(self : StreamingSession) -> Unit {
self.stats = OnlineStats::new()
self.intervals = []
self.started = false
self.finished = false
}
///|
/// Analyze a batch as a set of independent streaming sessions.
pub fn analyze_streams(
streams : Array[Array[Double]],
config : HrvConfig,
) -> Array[StreamingSnapshot] {
let result = []
for values in streams {
let session = StreamingSession::new(config)
session.push_all(values)
session.finish()
result.push(session.snapshot())
}
result
}
///|
/// Return the average RMSSD across non-empty snapshots.
pub fn average_stream_rmssd(snapshots : Array[StreamingSnapshot]) -> Double {
let values = []
for snapshot in snapshots {
if snapshot.sample_count > 0 {
values.push(snapshot.metrics.rmssd)
}
}
mean_value(values)
}