use observation::{Event, EventScope};
use std::collections::VecDeque;
use std::sync::{Arc, Mutex};
#[derive(Debug, Clone)]
pub struct RingConfig {
pub max_events: usize,
pub max_bytes: usize,
pub flush_bytes: usize,
pub flush_ms: u64,
}
impl Default for RingConfig {
fn default() -> Self {
Self {
max_events: 10_000,
max_bytes: 64 * 1024 * 1024, flush_bytes: 4 * 1024, flush_ms: 250,
}
}
}
struct Entry {
cursor: u64,
scope: EventScope,
event: Event,
byte_cost: usize,
}
struct Inner {
entries: VecDeque<Entry>,
next_cursor: u64,
total_bytes: usize,
pending_bytes: usize,
}
pub struct EventRing {
cfg: RingConfig,
inner: Mutex<Inner>,
}
impl EventRing {
pub fn new(cfg: RingConfig) -> Arc<Self> {
Arc::new(Self {
cfg,
inner: Mutex::new(Inner {
entries: VecDeque::new(),
next_cursor: 0,
total_bytes: 0,
pending_bytes: 0,
}),
})
}
pub fn push(&self, scope: EventScope, event: Event) -> (u64, bool) {
let byte_cost = event.fields.to_string().len();
let mut g = self.inner.lock().unwrap();
let cursor = g.next_cursor;
g.next_cursor += 1;
while (g.entries.len() >= self.cfg.max_events
|| g.total_bytes + byte_cost > self.cfg.max_bytes)
&& !g.entries.is_empty()
{
let evicted = g.entries.pop_front().unwrap();
g.total_bytes -= evicted.byte_cost;
}
g.total_bytes += byte_cost;
g.pending_bytes += byte_cost;
g.entries.push_back(Entry { cursor, scope, event, byte_cost });
let should_flush = g.pending_bytes >= self.cfg.flush_bytes;
(cursor, should_flush)
}
pub fn take_pending(&self) -> Vec<(EventScope, Event)> {
let mut g = self.inner.lock().unwrap();
g.pending_bytes = 0;
g.entries.iter().map(|e| (e.scope.clone(), e.event.clone())).collect()
}
pub fn tail_since(
&self,
scope: &EventScope,
since_cursor: u64,
limit: usize,
) -> (Vec<Event>, u64) {
let g = self.inner.lock().unwrap();
let mut last_cursor = since_cursor;
let events: Vec<Event> = g
.entries
.iter()
.filter(|e| e.cursor >= since_cursor && &e.scope == scope)
.take(limit)
.map(|e| {
last_cursor = e.cursor + 1;
e.event.clone()
})
.collect();
let next = if limit > 0 && events.len() == limit {
last_cursor
} else {
g.next_cursor
};
drop(g);
(events, next)
}
pub fn next_cursor(&self) -> u64 {
self.inner.lock().unwrap().next_cursor
}
pub fn len(&self) -> usize {
self.inner.lock().unwrap().entries.len()
}
pub fn is_empty(&self) -> bool {
self.len() == 0
}
}
#[cfg(test)]
mod tests {
use super::*;
use observation::{EventSource, Level, TaskRunId};
use serde_json::json;
fn make_event(seq: u32) -> Event {
Event {
run_id: TaskRunId::new(),
seq,
offset_ms: seq * 10,
level: Level::Info,
target: "test".to_string(),
msg: format!("event {seq}"),
fields: json!({}),
anchor: None,
source: EventSource::Synth,
}
}
fn make_scope() -> EventScope {
EventScope::TaskRun(TaskRunId::new())
}
#[test]
fn ring_basic_push_tail() {
let ring = EventRing::new(RingConfig::default());
let scope = make_scope();
for i in 0..5 {
ring.push(scope.clone(), make_event(i));
}
let (events, _next) = ring.tail_since(&scope, 0, 100);
assert_eq!(events.len(), 5);
}
#[test]
fn ring_evicts_at_capacity() {
let cfg = RingConfig { max_events: 3, max_bytes: usize::MAX, ..Default::default() };
let ring = EventRing::new(cfg);
let scope = make_scope();
for i in 0..5 {
ring.push(scope.clone(), make_event(i));
}
assert_eq!(ring.len(), 3);
}
#[test]
fn ring_no_drops_within_quota() {
let ring = EventRing::new(RingConfig::default());
let scope = make_scope();
for i in 0..10_000 {
ring.push(scope.clone(), make_event(i as u32));
}
assert_eq!(ring.len(), 10_000);
}
#[test]
fn ring_flush_threshold() {
let cfg = RingConfig { flush_bytes: 1, ..Default::default() };
let ring = EventRing::new(cfg);
let scope = make_scope();
let (_, should_flush) = ring.push(scope.clone(), make_event(0));
assert!(should_flush);
}
#[test]
fn ring_take_pending_resets() {
let ring = EventRing::new(RingConfig::default());
let scope = make_scope();
ring.push(scope.clone(), make_event(0));
let pending = ring.take_pending();
assert_eq!(pending.len(), 1);
let (_, should_flush) = ring.push(scope.clone(), make_event(1));
assert!(!should_flush);
}
}