use std::collections::HashMap;
use std::collections::VecDeque;
use std::sync::{Arc, Mutex};
use stateset_core::CommerceEvent;
use stateset_embedded::Commerce;
use tokio::sync::broadcast;
pub(crate) const DEFAULT_REPLAY_CAPACITY: usize = 1024;
const LIVE_CHANNEL_CAPACITY: usize = 1024;
pub(crate) const STREAM_RESET_EVENT: &str = "stream_reset";
#[derive(Clone, Debug)]
pub(crate) struct SequencedEvent {
pub(crate) id: u64,
pub(crate) event: CommerceEvent,
}
#[derive(Debug)]
struct Ring {
buf: VecDeque<SequencedEvent>,
capacity: usize,
last_id: u64,
}
impl Ring {
fn new(capacity: usize) -> Self {
Self {
buf: VecDeque::with_capacity(capacity.min(1024)),
capacity: capacity.max(1),
last_id: 0,
}
}
fn push(&mut self, event: SequencedEvent) {
self.last_id = event.id;
if self.buf.len() == self.capacity {
self.buf.pop_front();
}
self.buf.push_back(event);
}
fn oldest_id(&self) -> Option<u64> {
self.buf.front().map(|e| e.id)
}
}
#[derive(Debug)]
pub(crate) struct EventReplayBuffer {
ring: Mutex<Ring>,
live_tx: broadcast::Sender<SequencedEvent>,
}
pub(crate) struct ReplayPlan {
pub(crate) events: Vec<SequencedEvent>,
pub(crate) gap_detected: bool,
}
impl EventReplayBuffer {
fn new(capacity: usize) -> Self {
let (live_tx, _live_rx) = broadcast::channel(LIVE_CHANNEL_CAPACITY);
Self { ring: Mutex::new(Ring::new(capacity)), live_tx }
}
pub(crate) fn subscribe_live(&self) -> broadcast::Receiver<SequencedEvent> {
self.live_tx.subscribe()
}
fn record(&self, event: CommerceEvent) {
let sequenced = {
let mut ring = self.ring.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
let id = ring.last_id + 1;
let sequenced = SequencedEvent { id, event };
ring.push(sequenced.clone());
sequenced
};
let _ = self.live_tx.send(sequenced);
}
pub(crate) fn replay_after(&self, last_event_id: Option<u64>) -> ReplayPlan {
let Some(last) = last_event_id else {
return ReplayPlan { events: Vec::new(), gap_detected: false };
};
let ring = self.ring.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
let gap_detected = ring.oldest_id().is_some_and(|oldest| oldest > last + 1);
let events = ring.buf.iter().filter(|e| e.id > last).cloned().collect::<Vec<_>>();
ReplayPlan { events, gap_detected }
}
#[cfg(test)]
pub(crate) fn last_id(&self) -> u64 {
self.ring.lock().unwrap_or_else(std::sync::PoisonError::into_inner).last_id
}
}
#[derive(Clone, Debug, Default)]
pub(crate) struct EventReplayRegistry {
buffers: Arc<Mutex<HashMap<usize, Arc<EventReplayBuffer>>>>,
capacity: usize,
}
impl EventReplayRegistry {
pub(crate) fn new(capacity: usize) -> Self {
Self { buffers: Arc::new(Mutex::new(HashMap::new())), capacity: capacity.max(1) }
}
pub(crate) fn buffer_for(&self, commerce: &Arc<Commerce>) -> Arc<EventReplayBuffer> {
let key = Arc::as_ptr(commerce) as usize;
let mut buffers = self.buffers.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
if let Some(existing) = buffers.get(&key) {
return Arc::clone(existing);
}
let buffer = Arc::new(EventReplayBuffer::new(self.capacity));
spawn_pump(commerce, &buffer);
buffers.insert(key, Arc::clone(&buffer));
buffer
}
}
fn spawn_pump(commerce: &Arc<Commerce>, buffer: &Arc<EventReplayBuffer>) {
use tokio_stream::StreamExt as _;
let mut subscription = commerce.subscribe_events();
let buffer = Arc::clone(buffer);
tokio::spawn(async move {
while let Some(event) = subscription.next().await {
buffer.record(event);
}
});
}
#[cfg(test)]
mod tests {
use super::*;
use chrono::Utc;
use stateset_core::{CommerceEvent, CustomerId};
fn customer_event(email: &str) -> CommerceEvent {
CommerceEvent::CustomerCreated {
customer_id: CustomerId::new(),
email: email.to_string(),
timestamp: Utc::now(),
}
}
#[test]
fn ring_assigns_monotonic_ids() {
let buffer = EventReplayBuffer::new(8);
buffer.record(customer_event("a@example.com"));
buffer.record(customer_event("b@example.com"));
assert_eq!(buffer.last_id(), 2);
}
#[test]
fn replay_after_filters_by_id() {
let buffer = EventReplayBuffer::new(8);
for i in 0..5 {
buffer.record(customer_event(&format!("user{i}@example.com")));
}
let plan = buffer.replay_after(Some(2));
let ids: Vec<u64> = plan.events.iter().map(|e| e.id).collect();
assert_eq!(ids, vec![3, 4, 5]);
assert!(!plan.gap_detected);
}
#[test]
fn replay_after_none_is_empty() {
let buffer = EventReplayBuffer::new(8);
buffer.record(customer_event("a@example.com"));
let plan = buffer.replay_after(None);
assert!(plan.events.is_empty());
assert!(!plan.gap_detected);
}
#[test]
fn overflow_evicts_oldest_and_reports_gap() {
let buffer = EventReplayBuffer::new(3);
for i in 0..5 {
buffer.record(customer_event(&format!("user{i}@example.com")));
}
let plan = buffer.replay_after(Some(1));
assert!(plan.gap_detected, "evicted ids after last_event_id must report a gap");
let ids: Vec<u64> = plan.events.iter().map(|e| e.id).collect();
assert_eq!(ids, vec![3, 4, 5]);
}
#[test]
fn contiguous_replay_reports_no_gap() {
let buffer = EventReplayBuffer::new(3);
for i in 0..5 {
buffer.record(customer_event(&format!("user{i}@example.com")));
}
let plan = buffer.replay_after(Some(3));
assert!(!plan.gap_detected);
let ids: Vec<u64> = plan.events.iter().map(|e| e.id).collect();
assert_eq!(ids, vec![4, 5]);
}
#[test]
fn registry_returns_same_buffer_per_commerce() {
let commerce = Arc::new(Commerce::new(":memory:").expect("commerce"));
let registry = EventReplayRegistry::new(16);
let rt = tokio::runtime::Builder::new_current_thread().enable_all().build().unwrap();
let _guard = rt.enter();
let a = registry.buffer_for(&commerce);
let b = registry.buffer_for(&commerce);
assert!(Arc::ptr_eq(&a, &b));
}
}