pub mod error;
pub mod event;
pub use error::EventLogError;
pub use event::{Event, EventCfg};
use commonware_runtime::buffer::paged::CacheRef;
use commonware_storage::journal::contiguous::{Reader as _, variable};
use commonware_utils::{NZU16, NZU64, NZUsize};
use futures::StreamExt as _;
use std::num::{NonZeroU16, NonZeroU64, NonZeroUsize};
const REPLAY_BUFFER: NonZeroUsize = NZUsize!(1024);
#[derive(Debug, Clone)]
pub struct EventLogConfig {
pub partition: String,
pub items_per_section: NonZeroU64,
pub event_cfg: EventCfg,
pub page_size: NonZeroU16,
pub page_cache_pages: NonZeroUsize,
pub write_buffer: NonZeroUsize,
}
impl EventLogConfig {
#[must_use]
pub fn for_partition(partition: impl Into<String>) -> Self {
Self {
partition: partition.into(),
items_per_section: NZU64!(1024),
event_cfg: EventCfg::DEFAULT,
page_size: NZU16!(16384),
page_cache_pages: NZUsize!(64),
write_buffer: NZUsize!(65536),
}
}
}
pub struct EventLog<E>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler,
{
journal: variable::Journal<E, Event>,
}
impl<E> EventLog<E>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler,
{
pub async fn open(context: E, config: EventLogConfig) -> Result<Self, EventLogError> {
let page_cache = CacheRef::from_pooler(&context, config.page_size, config.page_cache_pages);
let journal_cfg = variable::Config {
partition: config.partition,
items_per_section: config.items_per_section,
compression: None,
codec_config: config.event_cfg,
page_cache,
write_buffer: config.write_buffer,
};
let journal = variable::Journal::init(context, journal_cfg).await?;
Ok(Self { journal })
}
pub async fn append(&self, event: &Event) -> Result<u64, EventLogError> {
Ok(self.journal.append(event).await?)
}
pub async fn len(&self) -> u64 {
self.journal.size().await
}
pub async fn is_empty(&self) -> bool {
self.len().await == 0
}
#[allow(clippy::significant_drop_tightening)]
pub async fn replay_with_positions(&self) -> Result<Vec<(u64, Event)>, EventLogError> {
let reader = self.journal.reader().await;
let start = reader.bounds().start;
let stream = reader.replay(REPLAY_BUFFER, start).await?;
futures::pin_mut!(stream);
let mut out = Vec::new();
while let Some(item) = stream.next().await {
out.push(item?);
}
Ok(out)
}
pub async fn replay(&self) -> Result<Vec<Event>, EventLogError> {
Ok(self
.replay_with_positions()
.await?
.into_iter()
.map(|(_pos, event)| event)
.collect())
}
pub async fn commit(&self) -> Result<(), EventLogError> {
Ok(self.journal.commit().await?)
}
pub async fn sync(&self) -> Result<(), EventLogError> {
Ok(self.journal.sync().await?)
}
}
#[cfg(test)]
mod tests {
use super::{Event, EventLog, EventLogConfig};
use commonware_runtime::{Runner, Supervisor as _, deterministic};
#[test]
fn doc_example_open_append_replay() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let log = EventLog::open(context, EventLogConfig::for_partition("conv-1"))
.await
.expect("open log");
log.append(&Event::new("user_msg", b"hello".to_vec()))
.await
.unwrap();
log.append(&Event::new("tool_call", b"\x01\x02".to_vec()))
.await
.unwrap();
log.commit().await.unwrap();
let events = log.replay().await.unwrap();
assert_eq!(events.len(), 2);
assert_eq!(events[0].kind, "user_msg");
});
}
#[test]
fn append_then_replay_preserves_order_and_payload() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let log = EventLog::open(context, EventLogConfig::for_partition("conv-order"))
.await
.expect("open");
let appended = vec![
Event::new("user_msg", b"what is 2+2?".to_vec()),
Event::new("planner_decision", vec![0xde, 0xad]),
Event::new("tool_call", vec![0x01, 0x02, 0x03]),
Event::new("tool_result", vec![0xff, 0x00, 0xff]),
];
for (i, event) in appended.iter().enumerate() {
let pos = log.append(event).await.expect("append");
assert_eq!(pos, i as u64, "positions are 0-indexed and contiguous");
}
log.commit().await.expect("commit");
assert_eq!(log.len().await, 4);
assert!(!log.is_empty().await);
let replayed = log.replay().await.expect("replay");
assert_eq!(replayed, appended);
let with_pos = log.replay_with_positions().await.expect("replay+pos");
let positions: Vec<u64> = with_pos.iter().map(|(p, _)| *p).collect();
assert_eq!(positions, vec![0, 1, 2, 3]);
assert_eq!(with_pos[2].1.payload, vec![0x01, 0x02, 0x03]);
});
}
#[test]
fn empty_log_replays_empty() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let log = EventLog::open(context, EventLogConfig::for_partition("conv-empty"))
.await
.expect("open");
assert!(log.is_empty().await);
assert_eq!(log.len().await, 0);
assert!(log.replay().await.expect("replay").is_empty());
});
}
#[test]
fn reopen_recovers_committed_events() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let cfg = EventLogConfig::for_partition("conv-reopen");
{
let log = EventLog::open(context.child("first"), cfg.clone())
.await
.expect("open first");
log.append(&Event::new("user_msg", b"persist me".to_vec()))
.await
.expect("append");
log.sync().await.expect("sync");
}
let log = EventLog::open(context.child("second"), cfg)
.await
.expect("reopen");
let replayed = log.replay().await.expect("replay");
assert_eq!(replayed.len(), 1);
assert_eq!(replayed[0].kind, "user_msg");
assert_eq!(replayed[0].payload, b"persist me".to_vec());
});
}
#[test]
fn deterministic_runs_are_reproducible() {
fn run() -> String {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let log =
EventLog::open(context.child("det"), EventLogConfig::for_partition("det"))
.await
.expect("open");
for i in 0..6u8 {
log.append(&Event::new(format!("kind-{i}"), vec![i; i as usize]))
.await
.expect("append");
}
log.commit().await.expect("commit");
let _ = log.replay().await.expect("replay");
context.auditor().state()
})
}
let first = run();
let second = run();
assert_eq!(first, second, "deterministic runtime must be reproducible");
}
}