pub mod checkpoint;
pub mod error;
pub mod event;
pub mod integrity;
mod metrics;
pub mod nav;
pub mod taint;
pub use checkpoint::EventCountCheckpoint;
pub use error::EventLogError;
pub use event::{Event, EventCfg};
pub use integrity::{
IntegrityError, MMR_SIGNED_ROOT_KIND, extend_and_sign, rebuild_from_events, verify_replay,
};
pub use taint::{
GrantedCapabilities, TrifectaLegs, TrustTag, any_untrusted, any_untrusted_excluding,
trifecta_legs,
};
pub fn init_metrics() {
metrics::force();
}
use commonware_runtime::buffer::paged::CacheRef;
use commonware_storage::journal::contiguous::{Contiguous as _, variable};
use commonware_utils::sync::AsyncMutex;
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 QuarantinedItem {
pub position: u64,
pub error: String,
}
#[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: AsyncMutex<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: AsyncMutex::new(journal),
})
}
pub async fn destroy(self) -> Result<(), EventLogError> {
Ok(self.journal.into_inner().destroy().await?)
}
pub async fn append(&self, event: &Event) -> Result<u64, EventLogError> {
let start = std::time::Instant::now(); let result = self.journal.lock().await.append(event).await;
metrics::record_append(result.is_ok(), start.elapsed());
Ok(result?)
}
pub async fn len(&self) -> u64 {
self.journal.lock().await.size()
}
pub async fn is_empty(&self) -> bool {
self.len().await == 0
}
pub async fn replay_with_positions(&self) -> Result<Vec<(u64, Event)>, EventLogError> {
let reader = self.journal.lock().await.snapshot().await?;
let start = reader.bounds().start;
let stream = reader.replay(start, REPLAY_BUFFER).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_from_with_positions(
&self,
start: u64,
) -> Result<Vec<(u64, Event)>, EventLogError> {
let reader = self.journal.lock().await.snapshot().await?;
let bounds = reader.bounds();
let from = start.max(bounds.start);
if from >= bounds.end {
return Ok(Vec::new());
}
let stream = reader.replay(from, REPLAY_BUFFER).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 replay_from(&self, start: u64) -> Result<Vec<Event>, EventLogError> {
Ok(self
.replay_from_with_positions(start)
.await?
.into_iter()
.map(|(_pos, event)| event)
.collect())
}
pub async fn replay_quarantining(
&self,
) -> Result<(Vec<(u64, Event)>, Vec<QuarantinedItem>), EventLogError> {
let reader = self.journal.lock().await.snapshot().await?;
let bounds = reader.bounds();
let mut ok = Vec::new();
let mut quarantined = Vec::new();
for position in bounds {
match reader.read(position).await {
Ok(event) => ok.push((position, event)),
Err(err) => quarantined.push(QuarantinedItem {
position,
error: err.to_string(),
}),
}
}
Ok((ok, quarantined))
}
pub async fn commit(&self) -> Result<(), EventLogError> {
Ok(self.journal.lock().await.commit().await?)
}
pub async fn sync(&self) -> Result<(), EventLogError> {
Ok(self.journal.lock().await.sync().await?)
}
}
#[cfg(test)]
mod tests {
use super::{Event, EventLog, EventLogConfig, EventLogError};
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 replay_from_offset_returns_tail_only() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let log = EventLog::open(context, EventLogConfig::for_partition("conv-from"))
.await
.expect("open");
for i in 0..5u8 {
log.append(&Event::new(format!("k{i}"), vec![i]))
.await
.expect("append");
}
log.commit().await.expect("commit");
assert_eq!(log.replay_with_positions().await.expect("replay").len(), 5);
let tail = log
.replay_from_with_positions(2)
.await
.expect("replay_from");
let positions: Vec<u64> = tail.iter().map(|(p, _)| *p).collect();
assert_eq!(positions, vec![2, 3, 4]);
assert_eq!(tail[0].1.kind, "k2");
assert_eq!(tail.last().expect("non-empty").1.kind, "k4");
assert!(log.replay_from(5).await.expect("from end").is_empty());
assert!(log.replay_from(99).await.expect("past end").is_empty());
assert_eq!(
log.replay_from(0).await.expect("from 0"),
log.replay().await.expect("replay")
);
});
}
#[test]
fn repair_replay_quarantining_reports_nothing_bad_on_a_healthy_log() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let log = EventLog::open(context, EventLogConfig::for_partition("conv-healthy"))
.await
.expect("open");
for i in 0..4u8 {
log.append(&Event::new(format!("k{i}"), vec![i]))
.await
.expect("append");
}
log.commit().await.expect("commit");
let (ok, quarantined) = log.replay_quarantining().await.expect("quarantine replay");
assert_eq!(ok.len(), 4);
assert!(quarantined.is_empty());
assert_eq!(ok, log.replay_with_positions().await.expect("replay"));
});
}
#[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 destroy_removes_the_partition_and_reopen_is_empty() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let cfg = EventLogConfig::for_partition("destroy-me");
let log = EventLog::open(context.child("first"), cfg.clone())
.await
.expect("open");
log.append(&Event::new("k", b"payload".to_vec()))
.await
.expect("append");
log.commit().await.expect("commit");
log.destroy().await.expect("destroy");
let log = EventLog::open(context.child("second"), cfg)
.await
.expect("reopen");
assert!(
log.replay().await.expect("replay").is_empty(),
"a destroyed partition must reopen empty"
);
let pos = log
.append(&Event::new("k2", b"fresh".to_vec()))
.await
.expect("append after destroy");
assert_eq!(pos, 0, "the fresh partition starts at position zero");
});
}
#[test]
fn reopen_recovers_synced_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 reopen_recovers_committed_events_after_commit_only() {
use commonware_runtime::{Runner as _, tokio as cw_tokio};
let dir =
std::env::temp_dir().join(format!("polyc-eventlog-commit-only-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let cfg = EventLogConfig::for_partition("conv-commit-only");
let write_cfg = cfg.clone();
let write_runner =
cw_tokio::Runner::new(cw_tokio::Config::default().with_storage_directory(dir.clone()));
write_runner.start(move |context| async move {
let log = EventLog::open(context, write_cfg).await.expect("open");
log.append(&Event::new("user_msg", b"one".to_vec()))
.await
.expect("append 0");
log.append(&Event::new("output_msg", b"two".to_vec()))
.await
.expect("append 1");
log.append(&Event::new("tool_call", b"three".to_vec()))
.await
.expect("append 2");
log.commit().await.expect("commit");
});
let read_cfg = cfg;
let read_runner =
cw_tokio::Runner::new(cw_tokio::Config::default().with_storage_directory(dir.clone()));
read_runner.start(move |context| async move {
let log = EventLog::open(context, read_cfg)
.await
.expect("reopen recovers committed-only data with no sync");
let replayed = log.replay().await.expect("replay");
assert_eq!(replayed.len(), 3, "all three committed events survive");
assert_eq!(replayed[0].kind, "user_msg");
assert_eq!(replayed[0].payload, b"one".to_vec());
assert_eq!(replayed[1].kind, "output_msg");
assert_eq!(replayed[1].payload, b"two".to_vec());
assert_eq!(replayed[2].kind, "tool_call");
assert_eq!(replayed[2].payload, b"three".to_vec());
let pos = log
.append(&Event::new("recovered", b"still writable".to_vec()))
.await
.expect("append continues after commit-only recovery");
assert_eq!(pos, 3, "the new append continues at position 3");
});
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn reopen_recovers_committed_events_after_commit_only_deterministic() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let cfg = EventLogConfig::for_partition("conv-reopen-commit-only");
{
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.commit().await.expect("commit");
}
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");
}
#[test]
fn repair_replay_quarantining_recovers_events_after_a_corrupted_earlier_section() {
use commonware_runtime::{Runner as _, tokio as cw_tokio};
let dir = std::env::temp_dir().join(format!(
"polyc-eventlog-repair-earlier-section-{}",
std::process::id()
));
let _ = std::fs::remove_dir_all(&dir);
let mut cfg = EventLogConfig::for_partition("conv-repair");
cfg.items_per_section = commonware_utils::NZU64!(1);
let write_cfg = cfg.clone();
let write_runner =
cw_tokio::Runner::new(cw_tokio::Config::default().with_storage_directory(dir.clone()));
write_runner.start(move |context| async move {
let log = EventLog::open(context, write_cfg).await.expect("open");
log.append(&Event::new("user_msg", vec![b'A'; 20]))
.await
.expect("append 0");
log.append(&Event::new("output_msg", vec![b'B'; 20]))
.await
.expect("append 1");
log.append(&Event::new("tool_call", vec![b'C'; 20]))
.await
.expect("append 2");
log.sync().await.expect("sync");
});
let data_file = dir.join("conv-repair_data").join("0000000000000000");
let mut bytes = std::fs::read(&data_file).expect("read section 0");
bytes[10] = 0xFF;
std::fs::write(&data_file, &bytes).expect("write corrupted section 0");
let read_cfg = cfg;
let read_runner =
cw_tokio::Runner::new(cw_tokio::Config::default().with_storage_directory(dir.clone()));
read_runner.start(move |context| async move {
let log = EventLog::open(context, read_cfg)
.await
.expect("reopen succeeds: only the final section is re-validated on open");
let plain_err = log
.replay()
.await
.expect_err("a corrupted item makes the whole streamed replay fail");
assert!(
matches!(plain_err, EventLogError::Journal(_)),
"unexpected error variant: {plain_err:?}"
);
let (ok, quarantined) = log
.replay_quarantining()
.await
.expect("quarantining replay reports positions, not an Err");
assert_eq!(quarantined.len(), 1, "exactly the corrupted item");
assert_eq!(quarantined[0].position, 0);
assert_eq!(
ok.len(),
2,
"the two events after the corrupted one recover"
);
assert_eq!(ok[0], (1, Event::new("output_msg", vec![b'B'; 20])));
assert_eq!(ok[1], (2, Event::new("tool_call", vec![b'C'; 20])));
});
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn torn_write_truncated_journal_reopens_and_replays() {
use commonware_runtime::{Runner as _, tokio as cw_tokio};
let dir =
std::env::temp_dir().join(format!("polyc-eventlog-torn-write-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let cfg = EventLogConfig::for_partition("conv-torn");
let write_cfg = cfg.clone();
let write_runner =
cw_tokio::Runner::new(cw_tokio::Config::default().with_storage_directory(dir.clone()));
write_runner.start(move |context| async move {
let log = EventLog::open(context, write_cfg).await.expect("open");
log.append(&Event::new("user_msg", vec![b'A'; 20]))
.await
.expect("append 0");
log.append(&Event::new("output_msg", vec![b'B'; 20]))
.await
.expect("append 1");
log.append(&Event::new("tool_call", vec![b'C'; 20]))
.await
.expect("append 2");
log.sync().await.expect("sync");
});
let data_file = dir.join("conv-torn_data").join("0000000000000000");
let full = std::fs::read(&data_file).expect("read section 0");
let marker = [b'C'; 20];
let payload_start = full
.windows(marker.len())
.position(|w| w == marker)
.expect("item 2's payload is present in the untruncated section");
let mut corrupted = full;
for b in &mut corrupted[payload_start + 10..] {
*b = 0;
}
std::fs::write(&data_file, &corrupted).expect("simulate a torn write");
let read_cfg = cfg;
let read_runner =
cw_tokio::Runner::new(cw_tokio::Config::default().with_storage_directory(dir.clone()));
read_runner.start(move |context| async move {
let log = EventLog::open(context, read_cfg)
.await
.expect("reopen recovers from a torn tail without erroring");
let events = log
.replay()
.await
.expect("replay succeeds (never hard-fails) after a torn write");
assert!(
events.len() <= 3,
"recovery must never fabricate events beyond what was appended"
);
for event in &events {
assert!(
[
Event::new("user_msg", vec![b'A'; 20]),
Event::new("output_msg", vec![b'B'; 20]),
]
.contains(event),
"recovery must never return the torn (never-fully-written) third item"
);
}
let pos = log
.append(&Event::new("recovered", b"still writable".to_vec()))
.await
.expect("the partition accepts appends again after torn-write recovery");
log.commit().await.expect("commit after recovery");
assert_eq!(
pos,
events.len() as u64,
"the new append continues from wherever recovery left off"
);
});
let _ = std::fs::remove_dir_all(&dir);
}
}