///|
/// Aggregation choices for a fixed streaming window.
pub(all) enum AggregationKind {
  MeanAggregate
  SumAggregate
  MinimumAggregate
  MaximumAggregate
  StandardDeviationAggregate
}

///|
pub struct WindowAggregate {
  start_timestamp : Int64
  end_timestamp : Int64
  count : Int
  value : Double
  summary : StatsSummary
}

///|
pub struct WindowAggregator {
  size : Int
  kind : AggregationKind
  window : DoubleWindow
  mut start_timestamp : Int64
  mut end_timestamp : Int64
  mut count : Int
}

///|
pub fn WindowAggregator::new(
  size : Int,
  kind? : AggregationKind = MeanAggregate,
) -> WindowAggregator {
  {
    size: if size < 1 {
      1
    } else {
      size
    },
    kind,
    window: DoubleWindow::new(if size < 1 { 1 } else { size }),
    start_timestamp: 0L,
    end_timestamp: 0L,
    count: 0,
  }
}

///|
fn WindowAggregator::value(self : WindowAggregator) -> Double {
  match self.kind {
    MeanAggregate => self.window.mean()
    SumAggregate => self.window.sum()
    MinimumAggregate => array_minimum(self.window.to_array())
    MaximumAggregate => array_maximum(self.window.to_array())
    StandardDeviationAggregate => self.window.standard_deviation()
  }
}

///|
fn WindowAggregator::close(self : WindowAggregator) -> WindowAggregate {
  {
    start_timestamp: self.start_timestamp,
    end_timestamp: self.end_timestamp,
    count: self.count,
    value: self.value(),
    summary: self.window.snapshot(),
  }
}

///|
pub fn WindowAggregator::push(
  self : WindowAggregator,
  point : SignalPoint,
) -> WindowAggregate? {
  if self.count == 0 {
    self.start_timestamp = point.timestamp
  }
  self.end_timestamp = point.timestamp
  self.count += 1
  ignore(self.window.push(point.value))
  if self.count >= self.size {
    let result = self.close()
    self.window.clear()
    self.count = 0
    Some(result)
  } else {
    None
  }
}

///|
pub fn WindowAggregator::flush(self : WindowAggregator) -> WindowAggregate? {
  if self.count == 0 {
    None
  } else {
    let result = self.close()
    self.window.clear()
    self.count = 0
    Some(result)
  }
}

///|
pub struct WatermarkTracker {
  allowed_lateness : Int64
  mut maximum_seen : Int64
  mut accepted : Int
  mut dropped : Int
}

///|
pub fn WatermarkTracker::new(
  allowed_lateness? : Int64 = 5L,
) -> WatermarkTracker {
  {
    allowed_lateness: if allowed_lateness < 0L {
      0L
    } else {
      allowed_lateness
    },
    maximum_seen: -9223372036854775807L,
    accepted: 0,
    dropped: 0,
  }
}

///|
pub fn WatermarkTracker::watermark(self : WatermarkTracker) -> Int64 {
  self.maximum_seen - self.allowed_lateness
}

///|
pub fn WatermarkTracker::observe(
  self : WatermarkTracker,
  timestamp : Int64,
) -> Bool {
  let accepted = self.maximum_seen == -9223372036854775807L ||
    timestamp >= self.watermark()
  if timestamp > self.maximum_seen {
    self.maximum_seen = timestamp
  }
  if accepted {
    self.accepted += 1
  } else {
    self.dropped += 1
  }
  accepted
}

///|
pub fn WatermarkTracker::accepted(self : WatermarkTracker) -> Int {
  self.accepted
}

///|
pub fn WatermarkTracker::dropped(self : WatermarkTracker) -> Int {
  self.dropped
}

///|
/// End-to-end stream state: reorder, aggregate, and detect.
pub struct StreamEngine {
  reorder : ReorderBuffer
  tracker : WatermarkTracker
  aggregator : WindowAggregator
  pipeline : MetricPipeline
  mut processed : Int
  mut aggregates : Int
  mut alerts : Int
}

///|
pub fn StreamEngine::new(
  metric : String,
  detector : PipelineDetector,
  window_size? : Int = 16,
  lateness? : Int = 4,
) -> StreamEngine {
  {
    reorder: ReorderBuffer::new(capacity=lateness + 1, policy=KeepForCorrection),
    tracker: WatermarkTracker::new(allowed_lateness=lateness.to_int64()),
    aggregator: WindowAggregator::new(window_size),
    pipeline: MetricPipeline::new(metric, detector),
    processed: 0,
    aggregates: 0,
    alerts: 0,
  }
}

///|
pub fn StreamEngine::process_point(
  self : StreamEngine,
  point : SignalPoint,
) -> Array[AlertEvent] {
  let result : Array[AlertEvent] = []
  let ready = self.reorder.push(point)
  for ordered in ready {
    if !self.tracker.observe(ordered.point.timestamp) {
      continue
    }
    self.processed += 1
    match self.aggregator.push(ordered.point) {
      None => ()
      Some(aggregate) => {
        self.aggregates += 1
        match self.pipeline.process(aggregate.end_timestamp, aggregate.value) {
          None => ()
          Some(alert) => {
            self.alerts += 1
            result.push(alert)
          }
        }
      }
    }
  }
  result
}

///|
pub fn StreamEngine::flush(self : StreamEngine) -> Array[AlertEvent] {
  let result : Array[AlertEvent] = []
  let pending = self.reorder.flush()
  for ordered in pending {
    if !self.tracker.observe(ordered.point.timestamp) {
      continue
    }
    self.processed += 1
    match self.aggregator.push(ordered.point) {
      None => ()
      Some(aggregate) => {
        self.aggregates += 1
        match self.pipeline.process(aggregate.end_timestamp, aggregate.value) {
          None => ()
          Some(alert) => {
            self.alerts += 1
            result.push(alert)
          }
        }
      }
    }
  }
  match self.aggregator.flush() {
    None => ()
    Some(aggregate) =>
      match self.pipeline.process(aggregate.end_timestamp, aggregate.value) {
        None => ()
        Some(alert) => result.push(alert)
      }
  }
  result
}

///|
pub fn StreamEngine::processed(self : StreamEngine) -> Int {
  self.processed
}

///|
pub fn StreamEngine::aggregate_count(self : StreamEngine) -> Int {
  self.aggregates
}

///|
pub fn StreamEngine::alert_count(self : StreamEngine) -> Int {
  self.alerts
}