///|
/// A bounded in-memory event buffer for adapters that receive transactions incrementally.
pub(all) struct EventBuffer {
  capacity : Int
  events : Array[Transaction]
  dropped : Int
}

///|
pub fn EventBuffer::new(capacity : Int) -> EventBuffer {
  { capacity: if capacity < 0 { 0 } else { capacity }, events: [], dropped: 0 }
}

///|
pub fn EventBuffer::push(self : EventBuffer, tx : Transaction) -> EventBuffer {
  let events : Array[Transaction] = []
  for old in self.events {
    events.push(old)
  }
  let mut dropped = self.dropped
  if self.capacity == 0 {
    dropped += 1
  } else {
    if events.length() >= self.capacity {
      let _ = events.remove(0)
      dropped += 1
    }
    events.push(tx)
  }
  { ..self, events, dropped }
}

///|
pub fn EventBuffer::size(self : EventBuffer) -> Int {
  self.events.length()
}

///|
pub fn EventBuffer::is_full(self : EventBuffer) -> Bool {
  self.capacity > 0 && self.size() >= self.capacity
}

///|
pub fn EventBuffer::snapshot(self : EventBuffer) -> Array[Transaction] {
  let result : Array[Transaction] = []
  for tx in self.events {
    result.push(tx)
  }
  result
}

///|
pub(all) struct StreamProcessor {
  rules : Array[Rule]
  buffer : EventBuffer
  alerts : Array[Alert]
}

///|
pub fn StreamProcessor::new(
  rules : Array[Rule],
  capacity : Int,
) -> StreamProcessor {
  { rules, buffer: EventBuffer::new(capacity), alerts: [] }
}

///|
pub fn StreamProcessor::ingest(
  self : StreamProcessor,
  tx : Transaction,
) -> StreamProcessor {
  let buffer = self.buffer.push(tx)
  let alerts = evaluate(self.rules, buffer.snapshot())
  { ..self, buffer, alerts: deduplicate_alerts(alerts) }
}

///|
pub fn StreamProcessor::drain_alerts(
  self : StreamProcessor,
) -> (StreamProcessor, Array[Alert]) {
  ({ ..self, alerts: [] }, self.alerts)
}

///|
pub fn StreamProcessor::recent(
  self : StreamProcessor,
  window : Int,
  now : Int,
) -> Array[Transaction] {
  transactions_in_window(
    self.buffer.snapshot(),
    TimeWindow::new(now - window, window),
  )
}