mirror of
https://github.com/herdrdev/herdr.git
synced 2026-09-24 16:01:13 +00:00
47 lines
1.2 KiB
Rust
47 lines
1.2 KiB
Rust
#[derive(Clone, Default)]
|
|
pub struct EventHub {
|
|
inner: std::sync::Arc<std::sync::Mutex<EventHubState>>,
|
|
}
|
|
|
|
#[derive(Default)]
|
|
struct EventHubState {
|
|
next_sequence: u64,
|
|
events: Vec<(u64, crate::api::schema::EventEnvelope)>,
|
|
}
|
|
|
|
impl EventHub {
|
|
const MAX_EVENTS: usize = 512;
|
|
|
|
pub fn push(&self, event: crate::api::schema::EventEnvelope) {
|
|
let Ok(mut state) = self.inner.lock() else {
|
|
return;
|
|
};
|
|
state.next_sequence += 1;
|
|
let sequence = state.next_sequence;
|
|
state.events.push((sequence, event));
|
|
let overflow = state.events.len().saturating_sub(Self::MAX_EVENTS);
|
|
if overflow > 0 {
|
|
state.events.drain(0..overflow);
|
|
}
|
|
}
|
|
|
|
pub fn events_after(&self, sequence: u64) -> Vec<(u64, crate::api::schema::EventEnvelope)> {
|
|
let Ok(state) = self.inner.lock() else {
|
|
return Vec::new();
|
|
};
|
|
state
|
|
.events
|
|
.iter()
|
|
.filter(|(event_sequence, _)| *event_sequence > sequence)
|
|
.cloned()
|
|
.collect()
|
|
}
|
|
|
|
pub fn current_sequence(&self) -> u64 {
|
|
let Ok(state) = self.inner.lock() else {
|
|
return 0;
|
|
};
|
|
state.next_sequence
|
|
}
|
|
}
|