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