pub mod checkpoint;
pub mod error;
mod metrics;
pub use checkpoint::EventCountCheckpoint;
pub use error::EventLogError;
pub use polyc_eventlog_model::integrity;
pub use polyc_eventlog_model::integrity::{
IntegrityError, MMR_SIGNED_ROOT_KIND, RootStanding, extend_and_sign, rebuild_from_events,
root_standing_with_trust, verify_extension_with_trust, verify_replay, verify_replay_with_trust,
};
pub use polyc_eventlog_model::nav;
pub use polyc_eventlog_model::taint;
pub use polyc_eventlog_model::taint::{
GrantedCapabilities, TrifectaLegs, TrustTag, any_untrusted, any_untrusted_excluding,
trifecta_legs,
};
pub use polyc_eventlog_model::{BoundedReplay, Event, EventCfg};
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,
{
config: EventLogConfig,
journal: AsyncMutex<variable::Journal<E, Event>>,
}
fn journal_config<E>(context: &E, config: &EventLogConfig) -> variable::Config<EventCfg>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler,
{
variable::Config {
partition: config.partition.clone(),
items_per_section: config.items_per_section,
compression: None,
codec_config: config.event_cfg,
page_cache: CacheRef::from_pooler(context, config.page_size, config.page_cache_pages),
write_buffer: config.write_buffer,
}
}
impl<E> EventLog<E>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler,
{
pub async fn open(context: E, config: EventLogConfig) -> Result<Self, EventLogError> {
let journal_cfg = journal_config(&context, &config);
let journal = variable::Journal::init(context, journal_cfg).await?;
Ok(Self {
config,
journal: AsyncMutex::new(journal),
})
}
pub async fn reset(self, context: E) -> Result<Self, EventLogError> {
let Self { config, journal } = self;
drop(journal.into_inner());
let journal_cfg = journal_config(&context, &config);
let journal = variable::Journal::init_at_size(context, journal_cfg, 0).await?;
Ok(Self {
config,
journal: AsyncMutex::new(journal),
})
}
pub async fn reclaim(self) -> Result<(), EventLogError> {
let events = self.len().await;
if events != 0 {
return Err(EventLogError::ReclaimNotEmpty { events });
}
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_with_positions_bounded(
&self,
max_bytes: u64,
) -> Result<BoundedReplay, 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 events = Vec::new();
let mut bytes_read: u64 = 0;
let mut budget_exceeded = false;
while let Some(item) = stream.next().await {
let (position, event) = item?;
bytes_read = bytes_read.saturating_add(event.payload.len() as u64);
events.push((position, event));
if bytes_read > max_bytes {
budget_exceeded = true;
break;
}
}
Ok(BoundedReplay {
events,
bytes_read,
budget_exceeded,
})
}
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_from_with_positions_bounded(
&self,
start: u64,
max_bytes: u64,
) -> Result<BoundedReplay, 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(BoundedReplay {
events: Vec::new(),
bytes_read: 0,
budget_exceeded: false,
});
}
let stream = reader.replay(from, REPLAY_BUFFER).await?;
futures::pin_mut!(stream);
let mut events = Vec::new();
let mut bytes_read: u64 = 0;
let mut budget_exceeded = false;
while let Some(item) = stream.next().await {
let (position, event) = item?;
bytes_read = bytes_read.saturating_add(event.payload.len() as u64);
events.push((position, event));
if bytes_read > max_bytes {
budget_exceeded = true;
break;
}
}
Ok(BoundedReplay {
events,
bytes_read,
budget_exceeded,
})
}
pub async fn replay_range_with_positions_bounded(
&self,
start: u64,
end: u64,
max_bytes: u64,
) -> Result<BoundedReplay, EventLogError> {
let reader = self.journal.lock().await.snapshot().await?;
let bounds = reader.bounds();
let from = start.max(bounds.start);
let until = end.min(bounds.end);
if from >= until {
return Ok(BoundedReplay {
events: Vec::new(),
bytes_read: 0,
budget_exceeded: false,
});
}
let stream = reader.replay(from, REPLAY_BUFFER).await?;
futures::pin_mut!(stream);
let mut events = Vec::new();
let mut bytes_read: u64 = 0;
let mut budget_exceeded = false;
while let Some(item) = stream.next().await {
let (position, event) = item?;
if position >= until {
break;
}
bytes_read = bytes_read.saturating_add(event.payload.len() as u64);
events.push((position, event));
if bytes_read > max_bytes {
budget_exceeded = true;
break;
}
}
Ok(BoundedReplay {
events,
bytes_read,
budget_exceeded,
})
}
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 replay_with_positions_bounded_stops_reading_mid_partition_once_the_budget_trips() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let log = EventLog::open(context, EventLogConfig::for_partition("conv-bounded"))
.await
.expect("open");
for i in 0..10u32 {
log.append(&Event::new(format!("k{i}"), vec![0u8; 1_000]))
.await
.expect("append");
}
log.commit().await.expect("commit");
let bounded = log
.replay_with_positions_bounded(3_500)
.await
.expect("bounded replay");
assert!(
bounded.budget_exceeded,
"3,500-byte budget over a 10,000-byte partition must trip"
);
assert_eq!(
bounded.events.len(),
4,
"replay must stop the instant cumulative bytes (4,000 after the 4th event) \
cross the 3,500 budget — reading a 5th event (or draining the whole partition) \
means the abort happened too late, or not at all"
);
assert_eq!(
bounded.bytes_read, 4_000,
"bytes_read must reflect exactly the events actually returned, not the whole \
partition's real 10,000 bytes"
);
assert!(
bounded.bytes_read < 10_000,
"peak materialized bytes must stay bounded well under the partition's real \
size — the whole point of stopping mid-stream"
);
let positions: Vec<u64> = bounded.events.iter().map(|(p, _)| *p).collect();
assert_eq!(positions, vec![0, 1, 2, 3]);
});
}
#[test]
fn replay_with_positions_bounded_reads_everything_when_under_budget() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let log = EventLog::open(context, EventLogConfig::for_partition("conv-under-budget"))
.await
.expect("open");
for i in 0..5u32 {
log.append(&Event::new(format!("k{i}"), vec![0u8; 100]))
.await
.expect("append");
}
log.commit().await.expect("commit");
let unbounded = log.replay_with_positions().await.expect("replay");
let bounded = log
.replay_with_positions_bounded(10_000)
.await
.expect("bounded replay");
assert!(
!bounded.budget_exceeded,
"500 bytes under a 10,000 byte budget must never trip"
);
assert_eq!(
bounded.events, unbounded,
"must match the unbounded replay exactly"
);
assert_eq!(bounded.bytes_read, 500);
});
}
#[test]
fn replay_from_with_positions_bounded_honors_both_start_and_the_byte_budget() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let log = EventLog::open(context, EventLogConfig::for_partition("conv-tail-bounded"))
.await
.expect("open");
for i in 0..10u32 {
log.append(&Event::new(format!("k{i}"), vec![0u8; 1_000]))
.await
.expect("append");
}
log.commit().await.expect("commit");
let bounded = log
.replay_from_with_positions_bounded(3, 2_500)
.await
.expect("bounded tail replay");
assert!(
bounded.budget_exceeded,
"a 2,500-byte budget over a 7,000-byte tail (positions 3..10) must trip"
);
let positions: Vec<u64> = bounded.events.iter().map(|(p, _)| *p).collect();
assert_eq!(
positions,
vec![3, 4, 5],
"must resume at position 3 (never re-reading 0..3) and stop on the event that \
CROSSES the 2,500 budget — that event is returned, matching \
`EventLog::replay_with_positions_bounded`'s own accumulate-push-then-check \
order, so the two differ only in where they start"
);
assert_eq!(
bounded.bytes_read, 3_000,
"bytes_read counts exactly the events actually returned, the budget-crossing \
one included — never the tail's whole 7,000 bytes"
);
});
}
#[test]
fn replay_range_excludes_the_end_position_exactly() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let log = EventLog::open(context, EventLogConfig::for_partition("conv-range-exact"))
.await
.expect("open");
for i in 0..10u32 {
log.append(&Event::new(format!("k{i}"), vec![0u8; 100]))
.await
.expect("append");
}
log.commit().await.expect("commit");
let bounded = log
.replay_range_with_positions_bounded(3, 6, u64::MAX)
.await
.expect("range replay");
let positions: Vec<u64> = bounded.events.iter().map(|(p, _)| *p).collect();
assert_eq!(
positions,
vec![3, 4, 5],
"[3, 6) is positions 3, 4, 5 — position 6 is outside the range and must not be \
returned"
);
assert!(
!bounded.budget_exceeded,
"the RANGE stopped this replay, not the budget; conflating the two would tell a \
caller its result was truncated when it is complete"
);
assert_eq!(
bounded.bytes_read, 300,
"the excluded end event must not spend the caller's byte budget either"
);
});
}
#[test]
fn replay_range_end_past_the_tail_clamps_to_the_tail() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let log = EventLog::open(context, EventLogConfig::for_partition("conv-range-clamp"))
.await
.expect("open");
for i in 0..4u32 {
log.append(&Event::new(format!("k{i}"), vec![0u8; 100]))
.await
.expect("append");
}
log.commit().await.expect("commit");
let bounded = log
.replay_range_with_positions_bounded(2, 9_999, u64::MAX)
.await
.expect("range replay");
let positions: Vec<u64> = bounded.events.iter().map(|(p, _)| *p).collect();
assert_eq!(positions, vec![2, 3], "clamped to the tail, not an error");
assert!(!bounded.budget_exceeded);
});
}
#[test]
fn replay_range_start_below_the_pruning_boundary_clamps_up_to_it() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let mut cfg = EventLogConfig::for_partition("conv-range-pruned");
cfg.items_per_section = commonware_utils::NZU64!(1);
let log = EventLog::open(context, cfg).await.expect("open");
for i in 0..6u32 {
log.append(&Event::new(format!("k{i}"), vec![0u8; 100]))
.await
.expect("append");
}
log.commit().await.expect("commit");
let pruned = log.journal.lock().await.prune(3).await.expect("prune");
assert!(
pruned,
"positions 0..3 must actually have been pruned — otherwise the clamp below is \
never exercised"
);
let bounded = log
.replay_range_with_positions_bounded(1, 5, u64::MAX)
.await
.expect("a start below the pruning boundary must clamp up to it, never error");
let positions: Vec<u64> = bounded.events.iter().map(|(p, _)| *p).collect();
assert_eq!(
positions,
vec![3, 4],
"clamped up to the pruning boundary (3), not the caller's stale start (1) — and \
still respecting the requested end (5)"
);
assert!(!bounded.budget_exceeded);
});
}
#[test]
fn replay_range_empty_and_inverted_are_empty_not_errors() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let log = EventLog::open(context, EventLogConfig::for_partition("conv-range-empty"))
.await
.expect("open");
for i in 0..5u32 {
log.append(&Event::new(format!("k{i}"), vec![0u8; 100]))
.await
.expect("append");
}
log.commit().await.expect("commit");
for (start, end, why) in [
(2u64, 2u64, "equal bounds are empty — the end is exclusive"),
(4, 1, "an inverted range is empty, never reordered"),
(99, 200, "a range entirely past the tail is empty"),
] {
let bounded = log
.replay_range_with_positions_bounded(start, end, u64::MAX)
.await
.expect("range replay");
assert!(bounded.events.is_empty(), "{why}");
assert_eq!(bounded.bytes_read, 0, "{why}");
assert!(
!bounded.budget_exceeded,
"an empty range must not report a tripped budget: {why}"
);
}
});
}
#[test]
fn replay_range_byte_cap_trips_inside_the_range() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let log = EventLog::open(context, EventLogConfig::for_partition("conv-range-budget"))
.await
.expect("open");
for i in 0..10u32 {
log.append(&Event::new(format!("k{i}"), vec![0u8; 1_000]))
.await
.expect("append");
}
log.commit().await.expect("commit");
let bounded = log
.replay_range_with_positions_bounded(2, 9, 2_500)
.await
.expect("range replay");
assert!(
bounded.budget_exceeded,
"a 2,500-byte budget over a 7,000-byte range must trip"
);
let positions: Vec<u64> = bounded.events.iter().map(|(p, _)| *p).collect();
assert_eq!(
positions,
vec![2, 3, 4],
"stops on the event that CROSSES the budget, which is returned — the same \
accumulate-push-then-check order the sibling bounded replays use, so the three \
differ only in where they start and stop"
);
assert_eq!(bounded.bytes_read, 3_000);
});
}
#[test]
fn replay_from_with_positions_bounded_past_the_end_is_empty_not_exceeded() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let log = EventLog::open(
context,
EventLogConfig::for_partition("conv-tail-bounded-empty"),
)
.await
.expect("open");
for i in 0..3u32 {
log.append(&Event::new(format!("k{i}"), vec![0u8; 1_000]))
.await
.expect("append");
}
log.commit().await.expect("commit");
let bounded = log
.replay_from_with_positions_bounded(99, 1)
.await
.expect("bounded tail replay past the end");
assert!(bounded.events.is_empty());
assert_eq!(bounded.bytes_read, 0);
assert!(
!bounded.budget_exceeded,
"an empty replay must never report the budget as exceeded"
);
});
}
#[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 reset_empties_the_partition_and_returns_a_usable_journal() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let cfg = EventLogConfig::for_partition("reset-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");
let log = log.reset(context.child("first")).await.expect("reset");
assert_eq!(log.len().await, 0, "the reset journal is empty");
assert!(log.replay().await.expect("replay").is_empty());
let pos = log
.append(&Event::new("k2", b"fresh".to_vec()))
.await
.expect("append after reset");
assert_eq!(pos, 0, "the reset journal restarts at position zero");
log.commit().await.expect("commit");
drop(log);
let log = EventLog::open(context.child("second"), cfg)
.await
.expect("reopen");
let replayed = log.replay().await.expect("replay");
assert_eq!(
replayed.len(),
1,
"the reset is durable: only the post-reset event survives"
);
assert_eq!(replayed[0].kind, "k2");
});
}
#[test]
fn reclaim_removes_the_partition_and_reopen_is_empty() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let cfg = EventLogConfig::for_partition("reclaim-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");
let log = log.reset(context.child("first")).await.expect("reset");
log.reclaim().await.expect("reclaim");
let log = EventLog::open(context.child("second"), cfg)
.await
.expect("reopen");
assert!(
log.replay().await.expect("replay").is_empty(),
"a reclaimed partition must reopen empty"
);
let pos = log
.append(&Event::new("k2", b"fresh".to_vec()))
.await
.expect("append after reclaim");
assert_eq!(pos, 0, "the fresh partition starts at position zero");
});
}
#[test]
fn reclaim_refuses_a_journal_that_still_holds_events() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let cfg = EventLogConfig::for_partition("reclaim-refuses");
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");
let error = log.reclaim().await.expect_err("reclaim must refuse");
assert!(
matches!(error, EventLogError::ReclaimNotEmpty { events: 1 }),
"expected a refusal naming the surviving event count, got {error}"
);
let log = EventLog::open(context.child("second"), cfg)
.await
.expect("reopen");
assert_eq!(
log.replay().await.expect("replay").len(),
1,
"the refused reclaim removed nothing"
);
});
}
#[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);
}
fn backing_blobs(
dir: &std::path::Path,
suffix: &str,
partition: &str,
) -> Vec<std::path::PathBuf> {
let mut found: Vec<std::path::PathBuf> =
std::fs::read_dir(dir.join(format!("{partition}{suffix}")))
.expect("partition storage directory")
.map(|entry| entry.expect("directory entry").path())
.filter(|path| path.is_file())
.collect();
found.sort();
found
}
fn seed_three_events(dir: &std::path::Path, cfg: &EventLogConfig) {
use commonware_runtime::{Runner as _, tokio as cw_tokio};
let cfg = cfg.clone();
let runner = cw_tokio::Runner::new(
cw_tokio::Config::default().with_storage_directory(dir.to_path_buf()),
);
runner.start(move |context| async move {
let log = EventLog::open(context, cfg).await.expect("open");
for i in 0..3u32 {
log.append(&Event::new(format!("k{i}"), vec![b'x'; 32]))
.await
.expect("append");
}
log.sync().await.expect("sync");
});
}
#[derive(Debug, PartialEq, Eq)]
enum Reopened {
Refused,
Holding(usize),
}
fn reopen_and_count(dir: &std::path::Path, cfg: &EventLogConfig) -> Reopened {
use commonware_runtime::{Runner as _, tokio as cw_tokio};
let cfg = cfg.clone();
let runner = cw_tokio::Runner::new(
cw_tokio::Config::default().with_storage_directory(dir.to_path_buf()),
);
runner.start(move |context| async move {
let Ok(log) = EventLog::open(context, cfg).await else {
return Reopened::Refused;
};
log.replay()
.await
.map_or(Reopened::Refused, |events| Reopened::Holding(events.len()))
})
}
#[test]
fn a_backing_blob_below_the_runtime_header_never_reopens_short() {
let root = std::env::temp_dir().join(format!(
"polyc-eventlog-partial-header-{}-{}",
std::process::id(),
"established"
));
let _ = std::fs::remove_dir_all(&root);
let mut observed = Vec::new();
for suffix in ["_data", "_offsets-blobs"] {
for raw_len in 0..8u64 {
let dir = root.join(format!("{}{raw_len}", suffix.trim_start_matches('_')));
std::fs::create_dir_all(&dir).expect("case directory");
let cfg = EventLogConfig::for_partition("conv-partial-header");
seed_three_events(&dir, &cfg);
let blobs = backing_blobs(&dir, suffix, "conv-partial-header");
assert!(
!blobs.is_empty(),
"a committed partition must leave a {suffix} blob to shorten"
);
for blob in &blobs {
let file = std::fs::OpenOptions::new()
.write(true)
.open(blob)
.expect("open the backing blob");
file.set_len(raw_len).expect("shorten below the header");
file.sync_all().expect("make the truncation durable");
}
let outcome = reopen_and_count(&dir, &cfg);
assert!(
matches!(outcome, Reopened::Refused | Reopened::Holding(0 | 3)),
"{suffix} shortened to {raw_len} byte(s) reopened SHORT as {outcome:?}; a \
partial history is the one answer no layer above can tell from a \
legitimately short conversation"
);
let expected = match suffix {
"_data" => Reopened::Holding(0),
_ => Reopened::Refused,
};
assert_eq!(
outcome, expected,
"{suffix} shortened to {raw_len} byte(s) answered {outcome:?}. This test \
pins what Commonware 2026.7.1 actually does, so a release that changes it \
fails here rather than silently moving what the layers above have to catch"
);
observed.push((suffix, raw_len, outcome));
let _ = std::fs::remove_dir_all(&dir);
}
}
let _ = std::fs::remove_dir_all(&root);
assert_eq!(
observed.len(),
16,
"eight prefix lengths for each of two blobs"
);
}
#[test]
fn an_interrupted_reset_survives_a_process_restart_without_a_partial_history() {
use commonware_runtime::{Runner as _, deterministic, deterministic::FaultConfig};
let cases = [
("open", FaultConfig::default().open(1.0)),
("write", FaultConfig::default().write(1.0)),
("sync", FaultConfig::default().sync(1.0)),
("resize", FaultConfig::default().resize(1.0)),
("remove", FaultConfig::default().remove(1.0)),
("scan", FaultConfig::default().scan(1.0)),
(
"torn write",
FaultConfig::default().write(1.0).partial_write(1.0),
),
(
"torn resize",
FaultConfig::default().resize(1.0).partial_resize(1.0),
),
];
let mut outcomes = Vec::new();
for (label, faults) in cases {
let partition = format!("conv-reset-crash-{}", label.replace(' ', "-"));
let seeded = partition.clone();
let runner = deterministic::Runner::timed(std::time::Duration::from_secs(30));
let ((), checkpoint) = runner.start_and_recover(move |context| async move {
let cfg = EventLogConfig::for_partition(seeded);
let log = EventLog::open(context.child("seed"), cfg.clone())
.await
.expect("open");
for i in 0..3u32 {
log.append(&Event::new(format!("k{i}"), vec![b'x'; 24]))
.await
.expect("append");
}
log.sync().await.expect("sync");
drop(log);
let log = EventLog::open(context.child("reset"), cfg)
.await
.expect("reopen before the reset");
*context.storage_fault_config().write() = faults;
drop(log.reset(context.child("reset")).await.ok());
*context.storage_fault_config().write() = FaultConfig::default();
});
let restarted = partition.clone();
let runner = deterministic::Runner::from(checkpoint);
let events = runner.start(move |context| async move {
let cfg = EventLogConfig::for_partition(restarted);
let log = EventLog::open(context.child("restart"), cfg)
.await
.expect("a reopen after an interrupted reset must mount");
log.replay().await.expect("replay").len()
});
assert!(
events == 0 || events == 3,
"a reset interrupted by a {label} failure left {events} event(s). A partial \
history is the outcome `EventLog::reset` exists to make impossible: no layer \
above can tell it from a legitimately short conversation"
);
outcomes.push((label, events));
}
assert!(
outcomes.iter().any(|&(_, events)| events == 3),
"no case stopped before the reset became durable: {outcomes:?}"
);
assert!(
outcomes.iter().any(|&(_, events)| events == 0),
"no case crossed the destructive boundary, so nothing here exercised recovery \
FROM it: {outcomes:?}"
);
}
#[test]
fn every_prefix_of_a_reclaim_still_reopens_as_an_empty_journal() {
use commonware_runtime::{Runner as _, tokio as cw_tokio};
let order: [(&str, bool); 7] = [
("_data", false),
("_data", true),
("_offsets-blobs", false),
("_offsets-blobs", true),
("_offsets-metadata", false),
("_offsets-metadata", false),
("_offsets-metadata", true),
];
for prefix in 0..=order.len() {
let dir = std::env::temp_dir().join(format!(
"polyc-eventlog-reclaim-prefix-{}-{prefix}",
std::process::id()
));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).expect("case directory");
let cfg = EventLogConfig::for_partition("conv-reclaim-prefix");
seed_three_events(&dir, &cfg);
let reset_runner = cw_tokio::Runner::new(
cw_tokio::Config::default().with_storage_directory(dir.clone()),
);
let reset_cfg = cfg.clone();
reset_runner.start(move |context| async move {
let log = EventLog::open(context.child("pre"), reset_cfg)
.await
.expect("open");
let log = log.reset(context.child("pre")).await.expect("reset");
assert_eq!(log.len().await, 0, "reclaim only ever runs on an empty log");
});
let mut removals = 0;
for (suffix, whole_partition) in order.iter().take(prefix) {
let path = dir.join(format!("conv-reclaim-prefix{suffix}"));
if !path.exists() {
continue;
}
if *whole_partition {
std::fs::remove_dir_all(&path).expect("remove the partition directory");
removals += 1;
} else if let Some(file) = std::fs::read_dir(&path)
.expect("partition directory")
.filter_map(Result::ok)
.map(|entry| entry.path())
.find(|path| path.is_file())
{
std::fs::remove_file(file).expect("unlink one blob");
removals += 1;
}
}
assert_eq!(
removals, prefix,
"prefix {prefix} was meant to remove {prefix} artifact(s) and removed {removals}"
);
let outcome = reopen_and_count(&dir, &cfg);
assert_eq!(
outcome,
Reopened::Holding(0),
"a reclaim interrupted after {prefix} of its {} removals answered {outcome:?}; \
the erase path reopens the partition before it removes anything, so a leftover \
that will not mount strands the erasure instead of finishing it",
order.len()
);
let _ = std::fs::remove_dir_all(&dir);
}
}
}