///|
/// 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)
}