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, extend_and_sign, rebuild_from_events,
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,
{
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_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 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);
}
}