use std::{collections::HashMap, path::PathBuf, sync::Arc};
use commonware_runtime::{Spawner as _, Supervisor as _};
use polyc_eventlog::{
BoundedReplay, Event, EventLog, EventLogConfig, EventLogError, QuarantinedItem,
};
use polyc_host::{Body, HostHandle, HostOptions};
use tokio::sync::{mpsc, oneshot};
use tokio_util::sync::CancellationToken;
use polyc_crypto::signing_role::JournalAttestationSigner;
use polyc_proto::kinds::{
COMPACTION_CHECKPOINT as BASE_COMPACTION_CHECKPOINT, parse as parse_kind,
};
mod metrics;
mod observer;
mod partition_name;
pub fn init_metrics() {
metrics::force();
}
pub use observer::{AppendNotification, EventLogObserver, MutationKind, MutationNotification};
pub use partition_name::{
EncodedPartitionName, MAX_ENCODED_PARTITION_BYTES, PartitionNameError, PartitionRef,
decode as decode_partition, encode as encode_partition,
};
#[cfg(feature = "test-util")]
pub mod test_util {
pub fn arm_root_append_failure(partition: &str) {
super::arm_root_append_failure(partition);
}
pub fn arm_append_failure(partition: &str, event_index: usize) {
super::arm_append_failure_at(partition, event_index);
}
}
pub struct CheckpointReplay {
pub checkpoint: Option<CheckpointFrame>,
pub tail: Vec<Event>,
}
pub struct CheckpointFrame {
pub tail_offset: u64,
pub payload: Vec<u8>,
}
pub enum RewriteDecision {
Keep,
Drop,
Replace(Vec<u8>),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum PartitionPresence {
Absent,
Present,
Damaged {
remnant: String,
},
}
const OFFSETS_SUFFIXES: [&str; 2] = ["_offsets-blobs", "_offsets-metadata"];
const CHECKPOINT_STORAGE_SUFFIX: &str = "__eventcount_checkpoint";
enum Command {
Replay {
partition: PartitionRef,
ack: oneshot::Sender<Result<Vec<Event>, EventLogError>>,
},
ReplayWithPositions {
partition: PartitionRef,
ack: oneshot::Sender<Result<Vec<(u64, Event)>, EventLogError>>,
},
ReplayWithPositionsBounded {
partition: PartitionRef,
max_bytes: u64,
ack: oneshot::Sender<Result<BoundedReplay, EventLogError>>,
},
ReplayTailFromCheckpoint {
partition: PartitionRef,
ack: oneshot::Sender<Result<CheckpointReplay, EventLogError>>,
},
ReplayFromWithPositions {
partition: PartitionRef,
start: u64,
ack: oneshot::Sender<Result<Vec<(u64, Event)>, EventLogError>>,
},
ReplayFromWithPositionsBounded {
partition: PartitionRef,
start: u64,
max_bytes: u64,
ack: oneshot::Sender<Result<BoundedReplay, EventLogError>>,
},
ReplayRangeWithPositionsBounded {
partition: PartitionRef,
start: u64,
end: u64,
max_bytes: u64,
ack: oneshot::Sender<Result<BoundedReplay, EventLogError>>,
},
EventsSinceCheckpoint {
partition: PartitionRef,
ack: oneshot::Sender<Result<u64, EventLogError>>,
},
PartitionEventCount {
partition: PartitionRef,
ack: oneshot::Sender<Result<u64, AppendError>>,
},
DestroyPartition {
partition: PartitionRef,
ack: oneshot::Sender<Result<(), EventLogError>>,
},
AppendBatch {
partition: PartitionRef,
events: Vec<Event>,
ack: oneshot::Sender<Result<Vec<u64>, AppendError>>,
},
EnsurePartitionAttested {
partition: PartitionRef,
ack: oneshot::Sender<Result<Option<Event>, AppendError>>,
},
ResumeRewritePartition {
partition: PartitionRef,
command_id: String,
ack: oneshot::Sender<Result<Option<usize>, AppendError>>,
},
RewritePartition {
partition: PartitionRef,
command_id: String,
#[allow(clippy::type_complexity)]
decide: Box<dyn Fn(&Event) -> RewriteDecision + Send + Sync>,
ack: oneshot::Sender<Result<usize, EventLogError>>,
},
InspectRewrite {
partition: PartitionRef,
command_id: String,
#[allow(clippy::type_complexity)]
decide: Box<dyn Fn(&Event) -> RewriteDecision + Send + Sync>,
ack: oneshot::Sender<Result<RewritePlan, EventLogError>>,
},
ApplyRewrite {
partition: PartitionRef,
command_id: String,
before: [u8; 32],
after: [u8; 32],
stage_binding: [u8; 32],
#[allow(clippy::type_complexity)]
decide: Box<dyn Fn(&Event) -> RewriteDecision + Send + Sync>,
ack: oneshot::Sender<Result<RewriteApplication, EventLogError>>,
},
ListPartitions {
max_entries: usize,
max_name_bytes: usize,
excluded: Option<String>,
ack: oneshot::Sender<Result<Vec<String>, ListPartitionsError>>,
},
PartitionLastModified {
partition: PartitionRef,
ack: oneshot::Sender<Option<u64>>,
},
ProbePartition {
partition: PartitionRef,
ack: oneshot::Sender<Result<PartitionPresence, ListPartitionsError>>,
},
MigratePartition {
from: PartitionRef,
to: PartitionRef,
ack: oneshot::Sender<Result<PartitionMigration, AppendError>>,
},
MigrateAppendInto {
partition: PartitionRef,
events: Vec<Event>,
from: PartitionRef,
ack: oneshot::Sender<Result<usize, EventLogError>>,
},
ReplayVerified {
partition: PartitionRef,
ack: oneshot::Sender<Result<Vec<Event>, VerifyError>>,
},
RepairPartition {
partition: PartitionRef,
ack: oneshot::Sender<Result<Vec<QuarantinedItem>, EventLogError>>,
},
InspectRepair {
partition: PartitionRef,
replacement: Option<RepairEventReplacement>,
ack: oneshot::Sender<Result<RepairPlan, EventLogError>>,
},
ApplyRepair {
partition: PartitionRef,
before: [u8; 32],
after: [u8; 32],
replacement: Option<RepairEventReplacement>,
ack: oneshot::Sender<Result<RepairApplication, EventLogError>>,
},
RepairSettled {
partition: PartitionRef,
before: [u8; 32],
after: [u8; 32],
ack: oneshot::Sender<Result<bool, AppendError>>,
},
DestroyRepairStage {
partition: PartitionRef,
before: [u8; 32],
ack: oneshot::Sender<Result<(), AppendError>>,
},
DestroyRewriteStage {
partition: PartitionRef,
command_id: String,
ack: oneshot::Sender<Result<(), AppendError>>,
},
}
#[derive(Debug, Clone)]
pub struct RepairPlan {
pub before: [u8; 32],
pub after: [u8; 32],
pub will_rewrite: bool,
pub readable: Vec<Event>,
pub quarantined: Vec<QuarantinedItem>,
}
#[derive(Debug, Clone)]
pub struct RepairEventReplacement {
kind: String,
payload: Vec<u8>,
}
impl RepairEventReplacement {
#[must_use]
pub fn new(kind: impl Into<String>, payload: Vec<u8>) -> Self {
Self {
kind: kind.into(),
payload,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RepairApplication {
Applied,
AlreadyApplied,
Changed,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct RewritePlan {
pub before: [u8; 32],
pub after: [u8; 32],
pub stage_binding: [u8; 32],
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RewriteApplication {
Applied {
dropped: usize,
},
AlreadyApplied,
Changed,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PartitionMigration {
Migrated(usize),
EmptyDestroyed,
CopiedSourceRemains(usize),
}
#[derive(Debug, thiserror::Error)]
pub enum VerifyError {
#[error(transparent)]
Journal(#[from] EventLogError),
#[error(
"replay returned {actual} event(s) but the durable checkpoint expects at least \
{expected}: either the storage layer silently truncated the active section (a torn \
write or corruption/tampering), or a rewrite was interrupted before it lowered the \
count — see docs/operations/eventlog-disaster-recovery.md"
)]
TruncatedReplay {
expected: u64,
actual: u64,
},
#[error(transparent)]
Integrity(#[from] polyc_eventlog::IntegrityError),
}
const NUM_SHARDS: usize = 4;
fn shard_index(partition: &str, num_shards: usize) -> usize {
use std::hash::{Hash as _, Hasher as _};
let mut hasher = std::collections::hash_map::DefaultHasher::new();
partition.hash(&mut hasher);
#[allow(clippy::cast_possible_truncation)]
let hash = hasher.finish() as usize;
hash % num_shards
}
pub struct EventLogHost {
shards: Vec<HostHandle<Command>>,
host_shutdown: CancellationToken,
journal_attestation_public_key: Vec<u8>,
journal_attestation_trust: polyc_crypto::signing_role::RoleTrustSet<
polyc_crypto::signing_role::JournalAttestationRole,
>,
write_observers: observer::WriteObservers,
}
const COMMAND_BACKLOG: usize = 256;
impl EventLogHost {
#[allow(clippy::needless_pass_by_value)]
pub fn spawn(
storage_dir: PathBuf,
shutdown: CancellationToken,
journal_attestation_signer: JournalAttestationSigner,
) -> Result<Self, HostError> {
let trust = polyc_crypto::signing_role::RoleTrustSet::current(&journal_attestation_signer);
Self::spawn_with_trust(storage_dir, shutdown, journal_attestation_signer, trust)
}
#[allow(clippy::needless_pass_by_value)]
pub fn spawn_with_trust(
storage_dir: PathBuf,
shutdown: CancellationToken,
journal_attestation_signer: JournalAttestationSigner,
journal_attestation_trust: polyc_crypto::signing_role::RoleTrustSet<
polyc_crypto::signing_role::JournalAttestationRole,
>,
) -> Result<Self, HostError> {
let signer = journal_attestation_signer;
let journal_attestation_public_key = signer.public_key_bytes();
if !journal_attestation_trust
.keys()
.iter()
.any(|identity| identity == &signer.identity())
{
return Err(HostError::ActiveSignerNotTrusted);
}
let host_shutdown = shutdown.child_token();
let channels: Vec<(mpsc::Sender<Command>, mpsc::Receiver<Command>)> = (0..NUM_SHARDS)
.map(|_| mpsc::channel::<Command>(COMMAND_BACKLOG))
.collect();
let peer_senders: Vec<mpsc::Sender<Command>> =
channels.iter().map(|(tx, _)| tx.clone()).collect();
let write_observers = observer::WriteObservers::new();
let shard_ctx = ShardContext {
storage_dir: storage_dir.clone(),
signer,
journal_attestation_trust: journal_attestation_trust.clone(),
peer_senders,
write_observers: write_observers.clone(),
};
let mut shards = Vec::with_capacity(NUM_SHARDS);
for (idx, (tx, rx)) in channels.into_iter().enumerate() {
let body = ShardRouter {
ctx: shard_ctx.clone(),
workers: PartitionWorkers::new(),
};
let thread_name: &'static str =
Box::leak(format!("eventlog-host-{idx}").into_boxed_str());
let opts = HostOptions {
thread_name,
storage_dir: storage_dir.clone(),
command_backlog: COMMAND_BACKLOG,
};
let spawned =
polyc_host::spawn_host_with_channel(opts, host_shutdown.clone(), body, tx, rx);
match spawned {
Ok(handle) => shards.push(handle),
Err(_err) => {
host_shutdown.cancel();
return Err(HostError::RuntimeStart);
}
}
}
Ok(Self {
shards,
host_shutdown,
journal_attestation_public_key,
journal_attestation_trust,
write_observers,
})
}
#[must_use]
pub fn journal_attestation_public_key(&self) -> Vec<u8> {
self.journal_attestation_public_key.clone()
}
#[must_use]
pub fn journal_attestation_trust(
&self,
) -> polyc_crypto::signing_role::RoleTrustSet<polyc_crypto::signing_role::JournalAttestationRole>
{
self.journal_attestation_trust.clone()
}
pub fn register_observer(&self, observer: Arc<dyn EventLogObserver>) {
self.write_observers.register(observer);
}
#[must_use]
pub fn mutation_epoch(&self, partition: &str) -> u64 {
self.write_observers.epoch(partition)
}
#[must_use]
pub fn is_shutting_down(&self) -> bool {
self.host_shutdown.is_cancelled()
}
fn shard_for(&self, partition: &str) -> &HostHandle<Command> {
&self.shards[shard_index(partition, self.shards.len())]
}
#[cfg(test)]
#[must_use]
pub const fn num_shards(&self) -> usize {
self.shards.len()
}
#[cfg(test)]
#[must_use]
pub fn shard_index_for(&self, partition: &str) -> usize {
shard_index(partition, self.shards.len())
}
pub async fn list_partitions(&self) -> Result<Vec<String>, AppendError> {
self.list_partitions_bounded(usize::MAX, usize::MAX).await
}
pub async fn list_partitions_bounded(
&self,
max_entries: usize,
max_name_bytes: usize,
) -> Result<Vec<String>, AppendError> {
self.list_partitions_bounded_excluding(max_entries, max_name_bytes, None)
.await
}
pub async fn list_partitions_bounded_excluding(
&self,
max_entries: usize,
max_name_bytes: usize,
excluded: Option<String>,
) -> Result<Vec<String>, AppendError> {
let (ack, ack_rx) = oneshot::channel();
self.shards[0]
.sender()
.ok_or(AppendError::Closed)?
.send(Command::ListPartitions {
max_entries,
max_name_bytes,
excluded,
ack,
})
.await
.map_err(|_| AppendError::Closed)?;
ack_rx
.await
.map_err(|_| AppendError::Closed)?
.map_err(AppendError::Listing)
}
pub async fn partition_last_modified_ms(
&self,
partition: String,
) -> Result<Option<u64>, AppendError> {
let partition = PartitionRef::new(&partition)?;
let (ack, ack_rx) = oneshot::channel();
self.shards[0]
.sender()
.ok_or(AppendError::Closed)?
.send(Command::PartitionLastModified { partition, ack })
.await
.map_err(|_| AppendError::Closed)?;
ack_rx.await.map_err(|_| AppendError::Closed)
}
pub async fn partition_presence(
&self,
partition: String,
) -> Result<PartitionPresence, AppendError> {
let partition = PartitionRef::new(&partition)?;
let (ack, ack_rx) = oneshot::channel();
self.shards[0]
.sender()
.ok_or(AppendError::Closed)?
.send(Command::ProbePartition { partition, ack })
.await
.map_err(|_| AppendError::Closed)?;
ack_rx
.await
.map_err(|_| AppendError::Closed)?
.map_err(AppendError::Listing)
}
pub async fn migrate_partition(
&self,
from: String,
to: String,
) -> Result<PartitionMigration, AppendError> {
let from = PartitionRef::new(&from)?;
let to = PartitionRef::new(&to)?;
let (ack, ack_rx) = oneshot::channel();
self.shard_for(from.logical())
.sender()
.ok_or(AppendError::Closed)?
.send(Command::MigratePartition { from, to, ack })
.await
.map_err(|_| AppendError::Closed)?;
ack_rx.await.map_err(|_| AppendError::Closed)?
}
pub async fn destroy_partition(&self, partition: String) -> Result<(), AppendError> {
let partition = PartitionRef::new(&partition)?;
let (ack, ack_rx) = oneshot::channel();
self.shard_for(partition.logical())
.sender()
.ok_or(AppendError::Closed)?
.send(Command::DestroyPartition { partition, ack })
.await
.map_err(|_| AppendError::Closed)?;
ack_rx
.await
.map_err(|_| AppendError::Closed)?
.map_err(AppendError::Log)
}
#[cfg(test)]
pub async fn rewrite_partition_for_test(
&self,
partition: String,
decide: Box<dyn Fn(&Event) -> RewriteDecision + Send + Sync>,
) -> Result<usize, AppendError> {
let command_id = format!("test-rewrite-{partition}");
self.rewrite_partition(partition, command_id, decide).await
}
pub async fn resume_rewrite_partition(
&self,
partition: String,
command_id: String,
) -> Result<Option<usize>, AppendError> {
let partition = PartitionRef::new(&partition)?;
let (ack, ack_rx) = oneshot::channel();
self.shard_for(partition.logical())
.sender()
.ok_or(AppendError::Closed)?
.send(Command::ResumeRewritePartition {
partition,
command_id,
ack,
})
.await
.map_err(|_| AppendError::Closed)?;
ack_rx.await.map_err(|_| AppendError::Closed)?
}
pub async fn rewrite_partition(
&self,
partition: String,
command_id: String,
decide: Box<dyn Fn(&Event) -> RewriteDecision + Send + Sync>,
) -> Result<usize, AppendError> {
let partition = PartitionRef::new(&partition)?;
let (ack, ack_rx) = oneshot::channel();
self.shard_for(partition.logical())
.sender()
.ok_or(AppendError::Closed)?
.send(Command::RewritePartition {
partition,
command_id,
decide,
ack,
})
.await
.map_err(|_| AppendError::Closed)?;
ack_rx
.await
.map_err(|_| AppendError::Closed)?
.map_err(AppendError::Log)
}
pub async fn inspect_rewrite_partition(
&self,
partition: String,
command_id: String,
decide: Box<dyn Fn(&Event) -> RewriteDecision + Send + Sync>,
) -> Result<RewritePlan, AppendError> {
let partition = PartitionRef::new(&partition)?;
let (ack, ack_rx) = oneshot::channel();
self.shard_for(partition.logical())
.sender()
.ok_or(AppendError::Closed)?
.send(Command::InspectRewrite {
partition,
command_id,
decide,
ack,
})
.await
.map_err(|_| AppendError::Closed)?;
ack_rx
.await
.map_err(|_| AppendError::Closed)?
.map_err(AppendError::Log)
}
#[allow(clippy::too_many_arguments)]
pub async fn apply_rewrite_partition(
&self,
partition: String,
command_id: String,
before: [u8; 32],
after: [u8; 32],
stage_binding: [u8; 32],
decide: Box<dyn Fn(&Event) -> RewriteDecision + Send + Sync>,
) -> Result<RewriteApplication, AppendError> {
let partition = PartitionRef::new(&partition)?;
let (ack, ack_rx) = oneshot::channel();
self.shard_for(partition.logical())
.sender()
.ok_or(AppendError::Closed)?
.send(Command::ApplyRewrite {
partition,
command_id,
before,
after,
stage_binding,
decide,
ack,
})
.await
.map_err(|_| AppendError::Closed)?;
ack_rx
.await
.map_err(|_| AppendError::Closed)?
.map_err(AppendError::Log)
}
pub async fn replay(&self, partition: String) -> Result<Vec<Event>, AppendError> {
let partition = PartitionRef::new(&partition)?;
let (ack, ack_rx) = oneshot::channel();
self.shard_for(partition.logical())
.sender()
.ok_or(AppendError::Closed)?
.send(Command::Replay { partition, ack })
.await
.map_err(|_| AppendError::Closed)?;
ack_rx
.await
.map_err(|_| AppendError::Closed)?
.map_err(AppendError::Log)
}
pub async fn replay_with_positions(
&self,
partition: String,
) -> Result<Vec<(u64, Event)>, AppendError> {
let partition = PartitionRef::new(&partition)?;
let (ack, ack_rx) = oneshot::channel();
self.shard_for(partition.logical())
.sender()
.ok_or(AppendError::Closed)?
.send(Command::ReplayWithPositions { partition, ack })
.await
.map_err(|_| AppendError::Closed)?;
ack_rx
.await
.map_err(|_| AppendError::Closed)?
.map_err(AppendError::Log)
}
pub async fn replay_with_positions_bounded(
&self,
partition: String,
max_bytes: u64,
) -> Result<BoundedReplay, AppendError> {
let partition = PartitionRef::new(&partition)?;
let (ack, ack_rx) = oneshot::channel();
self.shard_for(partition.logical())
.sender()
.ok_or(AppendError::Closed)?
.send(Command::ReplayWithPositionsBounded {
partition,
max_bytes,
ack,
})
.await
.map_err(|_| AppendError::Closed)?;
ack_rx
.await
.map_err(|_| AppendError::Closed)?
.map_err(AppendError::Log)
}
pub async fn replay_from_with_positions(
&self,
partition: String,
start: u64,
) -> Result<Vec<(u64, Event)>, AppendError> {
let partition = PartitionRef::new(&partition)?;
let (ack, ack_rx) = oneshot::channel();
self.shard_for(partition.logical())
.sender()
.ok_or(AppendError::Closed)?
.send(Command::ReplayFromWithPositions {
partition,
start,
ack,
})
.await
.map_err(|_| AppendError::Closed)?;
ack_rx
.await
.map_err(|_| AppendError::Closed)?
.map_err(AppendError::Log)
}
pub async fn replay_from_with_positions_bounded(
&self,
partition: String,
start: u64,
max_bytes: u64,
) -> Result<BoundedReplay, AppendError> {
let partition = PartitionRef::new(&partition)?;
let (ack, ack_rx) = oneshot::channel();
self.shard_for(partition.logical())
.sender()
.ok_or(AppendError::Closed)?
.send(Command::ReplayFromWithPositionsBounded {
partition,
start,
max_bytes,
ack,
})
.await
.map_err(|_| AppendError::Closed)?;
ack_rx
.await
.map_err(|_| AppendError::Closed)?
.map_err(AppendError::Log)
}
pub async fn replay_range_with_positions_bounded(
&self,
partition: String,
start: u64,
end: u64,
max_bytes: u64,
) -> Result<BoundedReplay, AppendError> {
let partition = PartitionRef::new(&partition)?;
let (ack, ack_rx) = oneshot::channel();
self.shard_for(partition.logical())
.sender()
.ok_or(AppendError::Closed)?
.send(Command::ReplayRangeWithPositionsBounded {
partition,
start,
end,
max_bytes,
ack,
})
.await
.map_err(|_| AppendError::Closed)?;
ack_rx
.await
.map_err(|_| AppendError::Closed)?
.map_err(AppendError::Log)
}
pub async fn replay_tail_from_checkpoint(
&self,
partition: String,
) -> Result<CheckpointReplay, AppendError> {
let partition = PartitionRef::new(&partition)?;
let (ack, ack_rx) = oneshot::channel();
self.shard_for(partition.logical())
.sender()
.ok_or(AppendError::Closed)?
.send(Command::ReplayTailFromCheckpoint { partition, ack })
.await
.map_err(|_| AppendError::Closed)?;
ack_rx
.await
.map_err(|_| AppendError::Closed)?
.map_err(AppendError::Log)
}
pub async fn events_since_checkpoint(&self, partition: String) -> Result<u64, AppendError> {
let partition = PartitionRef::new(&partition)?;
let (ack, ack_rx) = oneshot::channel();
self.shard_for(partition.logical())
.sender()
.ok_or(AppendError::Closed)?
.send(Command::EventsSinceCheckpoint { partition, ack })
.await
.map_err(|_| AppendError::Closed)?;
ack_rx
.await
.map_err(|_| AppendError::Closed)?
.map_err(AppendError::Log)
}
pub async fn partition_event_count(&self, partition: String) -> Result<u64, AppendError> {
let partition = PartitionRef::new(&partition)?;
let (ack, ack_rx) = oneshot::channel();
self.shard_for(partition.logical())
.sender()
.ok_or(AppendError::Closed)?
.send(Command::PartitionEventCount { partition, ack })
.await
.map_err(|_| AppendError::Closed)?;
ack_rx.await.map_err(|_| AppendError::Closed)?
}
pub async fn append_batch(
&self,
partition: String,
events: Vec<Event>,
) -> Result<Vec<u64>, AppendError> {
let partition = PartitionRef::new(&partition)?;
for event in &events {
if event.payload.len() > polyc_proto::events_decode::MAX_EVENT_PAYLOAD_BYTES {
return Err(AppendError::PayloadTooLarge {
kind: event.kind.clone(),
len: event.payload.len(),
});
}
}
let (ack, ack_rx) = oneshot::channel();
self.shard_for(partition.logical())
.sender()
.ok_or(AppendError::Closed)?
.send(Command::AppendBatch {
partition,
events,
ack,
})
.await
.map_err(|_| AppendError::Closed)?;
ack_rx.await.map_err(|_| AppendError::Closed)?
}
pub async fn ensure_partition_attested(
&self,
partition: String,
) -> Result<Option<Event>, AppendError> {
let partition = PartitionRef::new(&partition)?;
let (ack, ack_rx) = oneshot::channel();
self.shard_for(partition.logical())
.sender()
.ok_or(AppendError::Closed)?
.send(Command::EnsurePartitionAttested { partition, ack })
.await
.map_err(|_| AppendError::Closed)?;
ack_rx.await.map_err(|_| AppendError::Closed)?
}
pub async fn replay_verified(&self, partition: String) -> Result<Vec<Event>, AppendError> {
let partition = PartitionRef::new(&partition)?;
let (ack, ack_rx) = oneshot::channel();
self.shard_for(partition.logical())
.sender()
.ok_or(AppendError::Closed)?
.send(Command::ReplayVerified { partition, ack })
.await
.map_err(|_| AppendError::Closed)?;
ack_rx
.await
.map_err(|_| AppendError::Closed)?
.map_err(AppendError::Verify)
}
pub async fn repair_partition(
&self,
partition: String,
) -> Result<Vec<QuarantinedItem>, AppendError> {
let partition = PartitionRef::new(&partition)?;
let (ack, ack_rx) = oneshot::channel();
self.shard_for(partition.logical())
.sender()
.ok_or(AppendError::Closed)?
.send(Command::RepairPartition { partition, ack })
.await
.map_err(|_| AppendError::Closed)?;
ack_rx
.await
.map_err(|_| AppendError::Closed)?
.map_err(AppendError::Log)
}
pub async fn inspect_repair_partition(
&self,
partition: String,
replacement: Option<RepairEventReplacement>,
) -> Result<RepairPlan, AppendError> {
let partition = PartitionRef::new(&partition)?;
let (ack, ack_rx) = oneshot::channel();
self.shard_for(partition.logical())
.sender()
.ok_or(AppendError::Closed)?
.send(Command::InspectRepair {
partition,
replacement,
ack,
})
.await
.map_err(|_| AppendError::Closed)?;
ack_rx
.await
.map_err(|_| AppendError::Closed)?
.map_err(AppendError::Log)
}
pub async fn apply_repair_partition(
&self,
partition: String,
before: [u8; 32],
after: [u8; 32],
replacement: Option<RepairEventReplacement>,
) -> Result<RepairApplication, AppendError> {
let partition = PartitionRef::new(&partition)?;
let (ack, ack_rx) = oneshot::channel();
self.shard_for(partition.logical())
.sender()
.ok_or(AppendError::Closed)?
.send(Command::ApplyRepair {
partition,
before,
after,
replacement,
ack,
})
.await
.map_err(|_| AppendError::Closed)?;
ack_rx
.await
.map_err(|_| AppendError::Closed)?
.map_err(AppendError::Log)
}
pub async fn repair_settled(
&self,
partition: String,
before: [u8; 32],
after: [u8; 32],
) -> Result<bool, AppendError> {
let partition = PartitionRef::new(&partition)?;
let (ack, ack_rx) = oneshot::channel();
self.shard_for(partition.logical())
.sender()
.ok_or(AppendError::Closed)?
.send(Command::RepairSettled {
partition,
before,
after,
ack,
})
.await
.map_err(|_| AppendError::Closed)?;
ack_rx.await.map_err(|_| AppendError::Closed)?
}
pub async fn destroy_rewrite_stage(
&self,
partition: String,
command_id: String,
) -> Result<(), AppendError> {
let partition = PartitionRef::new(&partition)?;
let (ack, ack_rx) = oneshot::channel();
self.shard_for(partition.logical())
.sender()
.ok_or(AppendError::Closed)?
.send(Command::DestroyRewriteStage {
partition,
command_id,
ack,
})
.await
.map_err(|_| AppendError::Closed)?;
ack_rx.await.map_err(|_| AppendError::Closed)?
}
pub async fn destroy_repair_stage(
&self,
partition: String,
before: [u8; 32],
) -> Result<(), AppendError> {
let partition = PartitionRef::new(&partition)?;
let (ack, ack_rx) = oneshot::channel();
self.shard_for(partition.logical())
.sender()
.ok_or(AppendError::Closed)?
.send(Command::DestroyRepairStage {
partition,
before,
ack,
})
.await
.map_err(|_| AppendError::Closed)?;
ack_rx.await.map_err(|_| AppendError::Closed)?
}
}
impl Drop for EventLogHost {
fn drop(&mut self) {
self.host_shutdown.cancel();
self.shards.clear();
}
}
#[derive(Debug, thiserror::Error)]
pub enum AppendError {
#[error("event-log host is shut down")]
Closed,
#[error(transparent)]
Log(#[from] EventLogError),
#[error("event append outcome is unknown: {0}")]
AppendOutcomeUnknown(EventLogError),
#[error("signed-root completion outcome is unknown: {0}")]
AttestationOutcomeUnknown(EventLogError),
#[error(transparent)]
Listing(#[from] ListPartitionsError),
#[error(transparent)]
Verify(#[from] VerifyError),
#[error(transparent)]
PartitionName(#[from] PartitionNameError),
#[error("event `{kind}` payload is {len} bytes, over the decode cap")]
PayloadTooLarge {
kind: String,
len: usize,
},
}
#[derive(Debug, thiserror::Error)]
pub enum HostError {
#[error("event-log runtime thread failed to start")]
RuntimeStart,
#[error("active journal-attestation signer is absent from its trust history")]
ActiveSignerNotTrusted,
}
struct ShardRouter {
ctx: ShardContext,
workers: PartitionWorkers,
}
#[derive(Clone)]
struct ShardContext {
storage_dir: PathBuf,
signer: JournalAttestationSigner,
journal_attestation_trust: polyc_crypto::signing_role::RoleTrustSet<
polyc_crypto::signing_role::JournalAttestationRole,
>,
peer_senders: Vec<mpsc::Sender<Command>>,
write_observers: observer::WriteObservers,
}
#[async_trait::async_trait]
impl Body for ShardRouter {
type Cmd = Command;
async fn open(&mut self, _ctx: &commonware_runtime::tokio::Context) -> Result<(), String> {
Ok(())
}
async fn handle(&mut self, ctx: &commonware_runtime::tokio::Context, cmd: Command) {
if let Command::ListPartitions {
max_entries,
max_name_bytes,
excluded,
ack,
} = cmd
{
let _ = ack.send(list_partitions_bounded_in(
&self.ctx.storage_dir,
max_entries,
max_name_bytes,
excluded.as_deref(),
));
return;
}
if let Command::PartitionLastModified { partition, ack } = cmd {
let _ = ack.send(last_modified_ms_in(&self.ctx.storage_dir, &partition));
return;
}
if let Command::ProbePartition { partition, ack } = cmd {
let _ = ack.send(
partition_presence_in(ctx, &self.ctx.storage_dir, partition.storage_key()).await,
);
return;
}
let partition = partition_key(&cmd).to_owned();
let tx = self.workers.get_or_spawn(ctx, &self.ctx, &partition).await;
if tx.send(cmd).await.is_err() {
tracing::warn!(
partition = %partition,
"partition worker channel closed immediately after dispatch; command dropped"
);
}
}
async fn on_drain_complete(&mut self, _ctx: &commonware_runtime::tokio::Context) {
self.workers.retire_all().await;
}
}
fn partition_key(cmd: &Command) -> &str {
match cmd {
Command::Replay { partition, .. }
| Command::ReplayWithPositions { partition, .. }
| Command::ReplayWithPositionsBounded { partition, .. }
| Command::ReplayTailFromCheckpoint { partition, .. }
| Command::ReplayFromWithPositions { partition, .. }
| Command::ReplayFromWithPositionsBounded { partition, .. }
| Command::ReplayRangeWithPositionsBounded { partition, .. }
| Command::EventsSinceCheckpoint { partition, .. }
| Command::PartitionEventCount { partition, .. }
| Command::DestroyPartition { partition, .. }
| Command::AppendBatch { partition, .. }
| Command::EnsurePartitionAttested { partition, .. }
| Command::RewritePartition { partition, .. }
| Command::InspectRewrite { partition, .. }
| Command::ApplyRewrite { partition, .. }
| Command::ResumeRewritePartition { partition, .. }
| Command::ReplayVerified { partition, .. }
| Command::RepairPartition { partition, .. }
| Command::InspectRepair { partition, .. }
| Command::ApplyRepair { partition, .. }
| Command::RepairSettled { partition, .. }
| Command::DestroyRepairStage { partition, .. }
| Command::DestroyRewriteStage { partition, .. }
| Command::MigrateAppendInto { partition, .. } => partition.logical(),
Command::MigratePartition { from, .. } => from.logical(),
Command::ListPartitions { .. }
| Command::PartitionLastModified { .. }
| Command::ProbePartition { .. } => {
unreachable!(
"the storage-directory reads are not partition-keyed; each is handled \
before this is called"
)
}
}
}
const PARTITION_WORKERS_CAP: usize = LOGS_CAP;
struct PartitionWorker {
tx: mpsc::Sender<Command>,
task: commonware_runtime::utils::Handle<()>,
}
fn metric_label_body(partition: &str) -> String {
let sanitized: String = partition
.chars()
.map(|c| if c.is_ascii_alphanumeric() { c } else { '_' })
.collect();
if sanitized.starts_with(|c: char| c.is_ascii_alphabetic()) {
sanitized
} else {
format!("p{sanitized}")
}
}
struct PartitionWorkers {
workers: HashMap<String, PartitionWorker>,
order: std::collections::VecDeque<String>,
worker_labels: HashMap<String, &'static str>,
}
impl PartitionWorkers {
fn new() -> Self {
Self {
workers: HashMap::new(),
order: std::collections::VecDeque::new(),
worker_labels: HashMap::new(),
}
}
fn worker_label(&mut self, partition: &str) -> &'static str {
if let Some(label) = self.worker_labels.get(partition) {
return label;
}
let label: &'static str =
Box::leak(format!("{}_worker", metric_label_body(partition)).into_boxed_str());
self.worker_labels.insert(partition.to_owned(), label);
label
}
fn touch(&mut self, partition: &str) {
self.order.retain(|p| p != partition);
self.order.push_back(partition.to_owned());
}
async fn get_or_spawn(
&mut self,
ctx: &commonware_runtime::tokio::Context,
shard_ctx: &ShardContext,
partition: &str,
) -> mpsc::Sender<Command> {
if let Some(worker) = self.workers.get(partition) {
let tx = worker.tx.clone();
self.touch(partition);
return tx;
}
let (tx, rx) = mpsc::channel::<Command>(COMMAND_BACKLOG);
let label = self.worker_label(partition);
let worker_ctx = ctx.child(label);
let shard_ctx = shard_ctx.clone();
let task =
worker_ctx.spawn(move |worker_ctx| partition_worker_loop(worker_ctx, rx, shard_ctx));
self.workers.insert(
partition.to_owned(),
PartitionWorker {
tx: tx.clone(),
task,
},
);
self.touch(partition);
while self.workers.len() > PARTITION_WORKERS_CAP {
let Some(evict) = self.order.pop_front() else {
break;
};
if let Some(worker) = self.workers.remove(&evict) {
Self::retire(&evict, worker).await;
}
}
tx
}
async fn retire(partition: &str, worker: PartitionWorker) {
drop(worker.tx);
if let Err(err) = worker.task.await {
tracing::warn!(
partition,
error = %err,
"partition worker task ended with an error while retiring"
);
}
}
async fn retire_all(&mut self) {
self.order.clear();
for (partition, worker) in self.workers.drain() {
Self::retire(&partition, worker).await;
}
}
}
async fn partition_worker_loop(
ctx: commonware_runtime::tokio::Context,
mut rx: mpsc::Receiver<Command>,
shard_ctx: ShardContext,
) {
let mut logs = OpenLogs::new();
let mut checkpoints = HashMap::new();
let mut count_checkpoints = CountCheckpoints::new();
let mut mmr_logs = HashMap::new();
while let Some(cmd) = rx.recv().await {
handle_command(
&ctx,
&shard_ctx.storage_dir,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&shard_ctx.signer,
&shard_ctx.journal_attestation_trust,
&shard_ctx.peer_senders,
&shard_ctx.write_observers,
cmd,
)
.await;
}
for log in logs.values() {
if let Err(err) = log.sync().await {
tracing::warn!(error = %err, "partition worker: event-log sync on retirement failed");
}
}
}
const LOGS_CAP: usize = 256;
struct OpenLogs<E>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
map: std::collections::HashMap<String, EventLog<E>>,
order: std::collections::VecDeque<String>,
log_labels: HashMap<String, &'static str>,
}
impl<E> OpenLogs<E>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
fn new() -> Self {
Self {
map: std::collections::HashMap::new(),
order: std::collections::VecDeque::new(),
log_labels: HashMap::new(),
}
}
fn log_label(&mut self, partition: &str) -> &'static str {
if let Some(label) = self.log_labels.get(partition) {
return label;
}
let label: &'static str = Box::leak(metric_label_body(partition).into_boxed_str());
self.log_labels.insert(partition.to_owned(), label);
label
}
fn values(&self) -> impl Iterator<Item = &EventLog<E>> {
self.map.values()
}
fn contains(&self, partition: &str) -> bool {
self.map.contains_key(partition)
}
fn take(&mut self, partition: &str) -> Option<EventLog<E>> {
self.order.retain(|p| p != partition);
self.map.remove(partition)
}
fn touch(&mut self, partition: &str) {
self.order.retain(|p| p != partition);
self.order.push_back(partition.to_owned());
}
async fn get_or_open(
&mut self,
context: &E,
partition: &str,
) -> Result<&EventLog<E>, EventLogError> {
if self.map.contains_key(partition) {
self.touch(partition);
return Ok(self
.map
.get(partition)
.expect("present per the contains check"));
}
let label = self.log_label(partition);
let ctx = context.child(label);
let log = EventLog::open(ctx, EventLogConfig::for_partition(partition)).await?;
self.map.insert(partition.to_owned(), log);
self.order.push_back(partition.to_owned());
while self.map.len() > LOGS_CAP {
let Some(evict) = self.order.pop_front() else {
break;
};
if let Some(log) = self.map.remove(&evict) {
if let Err(err) = log.sync().await {
tracing::warn!(partition = %evict, error = %err, "LRU evict: sync failed");
}
tracing::debug!(partition = %evict, "LRU-evicted open eventlog");
}
}
Ok(self.map.get(partition).expect("just inserted"))
}
}
#[allow(clippy::too_many_arguments, clippy::too_many_lines)]
async fn handle_command<E>(
context: &E,
storage_dir: &std::path::Path,
logs: &mut OpenLogs<E>,
checkpoints: &mut HashMap<String, CheckpointFrame>,
count_checkpoints: &mut CountCheckpoints<E>,
mmr_logs: &mut HashMap<String, polyc_mmr::VerifiableLog>,
signer: &JournalAttestationSigner,
journal_attestation_trust: &polyc_crypto::signing_role::RoleTrustSet<
polyc_crypto::signing_role::JournalAttestationRole,
>,
peer_senders: &[mpsc::Sender<Command>],
write_observers: &observer::WriteObservers,
cmd: Command,
) where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
match cmd {
Command::Replay { partition, ack } => {
let result = replay_one(context, logs, partition.storage_key()).await;
let _ = ack.send(result);
}
Command::ReplayWithPositions { partition, ack } => {
let result = replay_with_positions_one(context, logs, partition.storage_key()).await;
let _ = ack.send(result);
}
Command::ReplayWithPositionsBounded {
partition,
max_bytes,
ack,
} => {
let result = replay_with_positions_bounded_one(
context,
logs,
partition.storage_key(),
max_bytes,
)
.await;
let _ = ack.send(result);
}
Command::ReplayTailFromCheckpoint { partition, ack } => {
let result = replay_tail_one(context, logs, checkpoints, partition.storage_key()).await;
let _ = ack.send(result);
}
Command::ReplayFromWithPositions {
partition,
start,
ack,
} => {
let result =
replay_from_with_positions_one(context, logs, partition.storage_key(), start).await;
let _ = ack.send(result);
}
Command::ReplayFromWithPositionsBounded {
partition,
start,
max_bytes,
ack,
} => {
let result = replay_from_with_positions_bounded_one(
context,
logs,
partition.storage_key(),
start,
max_bytes,
)
.await;
let _ = ack.send(result);
}
Command::ReplayRangeWithPositionsBounded {
partition,
start,
end,
max_bytes,
ack,
} => {
let result = replay_range_with_positions_bounded_one(
context,
logs,
partition.storage_key(),
start,
end,
max_bytes,
)
.await;
let _ = ack.send(result);
}
Command::EventsSinceCheckpoint { partition, ack } => {
let result =
events_since_checkpoint_one(context, logs, checkpoints, partition.storage_key())
.await;
let _ = ack.send(result);
}
Command::PartitionEventCount { partition, ack } => {
let result =
partition_event_count_one(context, storage_dir, logs, partition.storage_key())
.await;
let _ = ack.send(result);
}
Command::AppendBatch {
partition,
events,
ack,
} => {
let result = append_batch_one(
context,
logs,
checkpoints,
count_checkpoints,
mmr_logs,
signer,
partition.storage_key(),
&events,
)
.await;
if let Ok(positions) = &result {
write_observers.notify_append(partition.logical(), &events, positions);
}
let _ = ack.send(result);
}
Command::EnsurePartitionAttested { partition, ack } => {
let result = ensure_partition_attested_one(
context,
storage_dir,
logs,
count_checkpoints,
mmr_logs,
signer,
journal_attestation_trust,
partition.storage_key(),
)
.await;
let _ = ack.send(result);
}
Command::DestroyPartition { partition, ack } => {
let already_erased = matches!(
Box::pin(partition_presence_in(
context,
storage_dir,
partition.storage_key()
))
.await,
Ok(PartitionPresence::Absent)
);
let result = if already_erased {
Ok(())
} else {
let destroyed = destroy_one(
context,
logs,
checkpoints,
count_checkpoints,
mmr_logs,
partition.storage_key(),
DestroyCheckpoint::Clear,
)
.await;
let verified = match destroyed {
Err(error) => Err(error),
Ok(()) => {
if matches!(
Box::pin(partition_presence_in(
context,
storage_dir,
partition.storage_key()
))
.await,
Ok(PartitionPresence::Absent)
) {
Ok(())
} else {
tracing::error!(
partition = %partition.storage_key(),
runbook = "docs/operations/eventlog-disaster-recovery.md",
"erasure reported success while durable storage survived it; \
the partition is not erased and the command is retryable"
);
Err(EventLogError::EraseIncomplete)
}
}
};
if verified.is_ok() {
write_observers.notify_mutation_bumping_epoch(
partition.logical(),
&MutationKind::Destroyed,
);
}
verified
};
let _ = ack.send(result);
}
Command::ResumeRewritePartition {
partition,
command_id,
ack,
} => {
let stage = rewrite_stage_partition(partition.storage_key(), &command_id);
let result = resume_rewrite(
context,
storage_dir,
logs,
checkpoints,
count_checkpoints,
mmr_logs,
partition.storage_key(),
&stage,
)
.await
.map_err(AppendError::Log);
if let Ok(Some(dropped)) = &result {
write_observers.notify_mutation_bumping_epoch(
partition.logical(),
&MutationKind::Rewritten { dropped: *dropped },
);
}
let _ = ack.send(result);
}
Command::RewritePartition {
partition,
command_id,
decide,
ack,
} => {
let result = rewrite_one(
context,
storage_dir,
logs,
checkpoints,
count_checkpoints,
mmr_logs,
signer,
&command_id,
partition.storage_key(),
decide.as_ref(),
)
.await;
if let Ok(dropped) = &result {
write_observers.notify_mutation_bumping_epoch(
partition.logical(),
&MutationKind::Rewritten { dropped: *dropped },
);
}
let _ = ack.send(result);
}
Command::InspectRewrite {
partition,
command_id,
decide,
ack,
} => {
let result = async {
let all = {
let log = logs.get_or_open(context, partition.storage_key()).await?;
log.replay().await?
};
inspect_exact_rewrite(
partition.storage_key(),
&command_id,
&all,
signer,
decide.as_ref(),
)
.map(|(plan, _, _)| plan)
}
.await;
let _ = ack.send(result);
}
Command::ApplyRewrite {
partition,
command_id,
before,
after,
stage_binding,
decide,
ack,
} => {
let result = apply_exact_rewrite(
context,
logs,
checkpoints,
count_checkpoints,
mmr_logs,
signer,
partition.storage_key(),
&command_id,
RewritePlan {
before,
after,
stage_binding,
},
decide.as_ref(),
)
.await;
if let Ok(RewriteApplication::Applied { dropped }) = &result {
write_observers.notify_mutation_bumping_epoch(
partition.logical(),
&MutationKind::Rewritten { dropped: *dropped },
);
}
let _ = ack.send(result);
}
Command::ListPartitions {
max_entries,
max_name_bytes,
excluded,
ack,
} => {
let _ = ack.send(list_partitions_bounded_in(
storage_dir,
max_entries,
max_name_bytes,
excluded.as_deref(),
));
}
Command::PartitionLastModified { partition, ack } => {
let _ = ack.send(last_modified_ms_in(storage_dir, &partition));
}
Command::ProbePartition { partition, ack } => {
let _ = ack
.send(partition_presence_in(context, storage_dir, partition.storage_key()).await);
}
Command::MigratePartition { from, to, ack } => {
let result = migrate_one(
context,
logs,
checkpoints,
count_checkpoints,
mmr_logs,
peer_senders,
&from,
&to,
write_observers,
)
.await;
let _ = ack.send(result);
}
Command::MigrateAppendInto {
partition,
events,
from,
ack,
} => {
let result = migrate_append_into(
context,
logs,
count_checkpoints,
mmr_logs,
signer,
&partition,
&events,
write_observers,
&from,
)
.await;
let _ = ack.send(result);
}
Command::ReplayVerified { partition, ack } => {
let result = replay_verified_one(
context,
logs,
count_checkpoints,
journal_attestation_trust,
partition.storage_key(),
)
.await;
let _ = ack.send(result);
}
Command::RepairPartition { partition, ack } => {
let result = repair_one(
context,
logs,
checkpoints,
count_checkpoints,
mmr_logs,
signer,
journal_attestation_trust,
partition.storage_key(),
)
.await;
if let Ok(effect) = &result
&& effect.rewrote
{
write_observers.notify_mutation_bumping_epoch(
partition.logical(),
&MutationKind::Repaired {
quarantined: effect.quarantined.len(),
},
);
}
let _ = ack.send(result.map(|effect| effect.quarantined));
}
Command::InspectRepair {
partition,
replacement,
ack,
} => {
let result = inspect_repair_one(
context,
logs,
count_checkpoints,
signer,
journal_attestation_trust,
partition.storage_key(),
replacement.as_ref(),
)
.await
.map(|scan| scan.plan());
let _ = ack.send(result);
}
Command::ApplyRepair {
partition,
before,
after,
replacement,
ack,
} => {
let result = apply_repair_conditionally(
context,
logs,
checkpoints,
count_checkpoints,
mmr_logs,
signer,
journal_attestation_trust,
partition.storage_key(),
RepairFingerprints { before, after },
replacement.as_ref(),
)
.await;
if let Ok((RepairApplication::Applied, quarantined)) = &result {
write_observers.notify_mutation_bumping_epoch(
partition.logical(),
&MutationKind::Repaired {
quarantined: *quarantined,
},
);
}
let _ = ack.send(result.map(|(application, _)| application));
}
Command::DestroyRewriteStage {
partition,
command_id,
ack,
} => {
let stage = rewrite_stage_partition(partition.storage_key(), &command_id);
let result = destroy_stage_if_present(
context,
storage_dir,
logs,
checkpoints,
count_checkpoints,
mmr_logs,
&stage,
)
.await;
let _ = ack.send(result);
}
Command::DestroyRepairStage {
partition,
before,
ack,
} => {
let stage = repair_stage_partition(partition.storage_key(), &before);
let result = destroy_stage_if_present(
context,
storage_dir,
logs,
checkpoints,
count_checkpoints,
mmr_logs,
&stage,
)
.await;
let _ = ack.send(result);
}
Command::RepairSettled {
partition,
before,
after,
ack,
} => {
let result = repair_settled_one(
context,
storage_dir,
logs,
count_checkpoints,
signer,
journal_attestation_trust,
partition.storage_key(),
RepairFingerprints { before, after },
)
.await;
let _ = ack.send(result);
}
}
}
struct CountCheckpoints<E>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
map: HashMap<String, polyc_eventlog::EventCountCheckpoint<E>>,
labels: HashMap<String, &'static str>,
}
impl<E> CountCheckpoints<E>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
fn new() -> Self {
Self {
map: HashMap::new(),
labels: HashMap::new(),
}
}
fn label(&mut self, partition: &str) -> &'static str {
if let Some(label) = self.labels.get(partition) {
return label;
}
let sanitized: String = partition
.chars()
.map(|c| if c.is_ascii_alphanumeric() { c } else { '_' })
.collect();
let label: &'static str = Box::leak(format!("{sanitized}_ckpt").into_boxed_str());
self.labels.insert(partition.to_owned(), label);
label
}
#[cfg(test)]
fn contains(&self, partition: &str) -> bool {
self.map.contains_key(partition)
}
fn forget(&mut self, partition: &str) {
self.map.remove(partition);
}
async fn take(
&mut self,
context: &E,
partition: &str,
) -> Result<polyc_eventlog::EventCountCheckpoint<E>, EventLogError> {
if let Some(checkpoint) = self.map.remove(partition) {
return Ok(checkpoint);
}
let label = self.label(partition);
Box::pin(polyc_eventlog::EventCountCheckpoint::open(
context.child(label),
partition,
))
.await
}
fn keep(&mut self, partition: &str, checkpoint: polyc_eventlog::EventCountCheckpoint<E>) {
self.map.insert(partition.to_owned(), checkpoint);
}
async fn expected_count(&mut self, context: &E, partition: &str) -> Result<u64, EventLogError> {
let checkpoint = self.take(context, partition).await?;
let count = checkpoint.expected_count();
self.keep(partition, checkpoint);
Ok(count)
}
async fn record(
&mut self,
context: &E,
partition: &str,
count: u64,
) -> Result<(), EventLogError> {
let checkpoint = self.take(context, partition).await?;
let checkpoint = Box::pin(checkpoint.record(count)).await?;
self.keep(partition, checkpoint);
Ok(())
}
async fn clear(&mut self, context: &E, partition: &str) -> Result<(), EventLogError> {
let checkpoint = self.take(context, partition).await?;
let checkpoint = Box::pin(checkpoint.clear()).await?;
self.keep(partition, checkpoint);
Ok(())
}
#[cfg(test)]
fn evict_all(&mut self) {
self.map.clear();
}
}
#[allow(clippy::too_many_arguments)]
async fn migrate_one<E>(
context: &E,
logs: &mut OpenLogs<E>,
checkpoints: &mut HashMap<String, CheckpointFrame>,
count_checkpoints: &mut CountCheckpoints<E>,
mmr_logs: &mut HashMap<String, polyc_mmr::VerifiableLog>,
peer_senders: &[mpsc::Sender<Command>],
from: &PartitionRef,
to: &PartitionRef,
write_observers: &observer::WriteObservers,
) -> Result<PartitionMigration, AppendError>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
if from == to {
return Ok(PartitionMigration::Migrated(0));
}
let source = {
let log = logs.get_or_open(context, from.storage_key()).await?;
log.replay().await?
};
if source.is_empty() {
destroy_one(
context,
logs,
checkpoints,
count_checkpoints,
mmr_logs,
from.storage_key(),
DestroyCheckpoint::Clear,
)
.await?;
write_observers.notify_mutation_bumping_epoch(
from.logical(),
&MutationKind::MigrationSource {
to: to.logical().to_owned(),
},
);
return Ok(PartitionMigration::EmptyDestroyed);
}
let to_shard = shard_index(to.logical(), peer_senders.len());
let carried: Vec<Event> = source
.into_iter()
.filter(|ev| ev.kind != polyc_eventlog::MMR_SIGNED_ROOT_KIND)
.collect();
let appended = if carried.is_empty() {
0
} else {
remote_migrate_append_into(peer_senders, to_shard, to, carried, from).await?
};
match destroy_one(
context,
logs,
checkpoints,
count_checkpoints,
mmr_logs,
from.storage_key(),
DestroyCheckpoint::Clear,
)
.await
{
Ok(()) => {
write_observers.notify_mutation_bumping_epoch(
from.logical(),
&MutationKind::MigrationSource {
to: to.logical().to_owned(),
},
);
Ok(PartitionMigration::Migrated(appended))
}
Err(err) => {
tracing::warn!(%from, %to, error = %err, "partition migration: events copied but the source destroy failed; a re-run retries the destroy without re-copying");
Ok(PartitionMigration::CopiedSourceRemains(appended))
}
}
}
#[allow(clippy::too_many_arguments)]
async fn migrate_append_into<E>(
context: &E,
logs: &mut OpenLogs<E>,
count_checkpoints: &mut CountCheckpoints<E>,
mmr_logs: &mut HashMap<String, polyc_mmr::VerifiableLog>,
signer: &JournalAttestationSigner,
to: &PartitionRef,
events: &[Event],
write_observers: &observer::WriteObservers,
from: &PartitionRef,
) -> Result<usize, EventLogError>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
if events.is_empty() {
return Ok(0);
}
let log = logs.get_or_open(context, to.storage_key()).await?;
let existing = log.replay().await?;
let present: std::collections::HashSet<(&str, &[u8])> = existing
.iter()
.map(|ev| (ev.kind.as_str(), ev.payload.as_slice()))
.collect();
let fresh: Vec<Event> = events
.iter()
.filter(|ev| !present.contains(&(ev.kind.as_str(), ev.payload.as_slice())))
.cloned()
.collect();
let tail_is_attested = existing
.last()
.is_some_and(|ev| ev.kind == polyc_eventlog::MMR_SIGNED_ROOT_KIND);
if fresh.is_empty() && tail_is_attested {
return Ok(0);
}
let appended = fresh.len();
let mut settled = existing;
settled.extend(fresh.iter().cloned());
let tree = polyc_eventlog::rebuild_from_events(&settled)?;
let marker = polyc_eventlog::extend_and_sign(&tree, &[], signer)?;
let mut to_write: Vec<Event> = fresh;
to_write.push(marker);
for (idx, event) in to_write.iter().enumerate() {
let append_result = match maybe_inject_append_failure(to.storage_key(), idx) {
Some(err) => Err(err),
None => log.append(event).await,
};
if let Err(err) = append_result {
if let Err(commit_err) = log.commit().await {
tracing::warn!(error = %commit_err, "commit of partial migration batch failed");
}
logs.take(to.storage_key());
mmr_logs.remove(to.storage_key());
return Err(err);
}
}
let commit_result = match maybe_inject_commit_failure(to.storage_key()) {
Some(err) => Err(err),
None => log.commit().await,
};
if let Err(err) = commit_result {
logs.take(to.storage_key());
mmr_logs.remove(to.storage_key());
return Err(err);
}
mmr_logs.remove(to.storage_key());
let new_len = log.len().await;
if let Err(err) = count_checkpoints
.record(context, to.storage_key(), new_len)
.await
{
tracing::warn!(%to, error = %err, "event-count checkpoint record failed after migration (#799)");
}
write_observers.notify_mutation_bumping_epoch(
to.logical(),
&MutationKind::MigrationDestination {
from: from.logical().to_owned(),
appended,
},
);
Ok(appended)
}
async fn remote_migrate_append_into(
peer_senders: &[mpsc::Sender<Command>],
shard: usize,
partition: &PartitionRef,
events: Vec<Event>,
from: &PartitionRef,
) -> Result<usize, AppendError> {
let (ack, ack_rx) = oneshot::channel();
peer_senders[shard]
.send(Command::MigrateAppendInto {
partition: partition.clone(),
events,
from: from.clone(),
ack,
})
.await
.map_err(|_| AppendError::Closed)?;
ack_rx
.await
.map_err(|_| AppendError::Closed)?
.map_err(AppendError::Log)
}
#[cfg(test)]
fn list_partitions_in(storage_dir: &std::path::Path) -> Result<Vec<String>, ListPartitionsError> {
list_partitions_bounded_in(storage_dir, usize::MAX, usize::MAX, None)
}
fn list_partitions_bounded_in(
storage_dir: &std::path::Path,
max_entries: usize,
max_name_bytes: usize,
excluded: Option<&str>,
) -> Result<Vec<String>, ListPartitionsError> {
const DATA_SUFFIX: &str = "_data";
let entries = match std::fs::read_dir(storage_dir) {
Ok(entries) => entries,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()),
Err(e) => return Err(ListPartitionsError::Storage(e)),
};
let mut partitions = Vec::new();
let mut name_bytes = 0_usize;
let mut data_dirs: std::collections::BTreeSet<String> = std::collections::BTreeSet::new();
let mut remnants: Vec<String> = Vec::new();
for entry in entries {
let entry = entry.map_err(ListPartitionsError::Storage)?;
if !entry
.file_type()
.map_err(ListPartitionsError::Storage)?
.is_dir()
{
continue;
}
let name = entry.file_name();
let Some(physical) = name.as_encoded_bytes().strip_suffix(DATA_SUFFIX.as_bytes()) else {
if let Some(remnant) = partition_shaped_remnant(&name) {
remnants.push(remnant);
}
continue;
};
let Some(physical) = std::str::from_utf8(physical).ok() else {
return Err(ListPartitionsError::Corrupt {
entry: name.to_string_lossy().into_owned(),
});
};
data_dirs.insert(physical.to_owned());
if physical.starts_with(REPAIR_STAGE_PREFIX)
|| physical.starts_with(REWRITE_STAGE_PREFIX)
|| physical.starts_with(EXACT_REWRITE_STAGE_PREFIX)
{
continue;
}
let logical =
partition_name::decode(physical).map_err(|_| ListPartitionsError::Corrupt {
entry: name.to_string_lossy().into_owned(),
})?;
if excluded == Some(logical.as_str()) {
continue;
}
let requested_entries =
partitions
.len()
.checked_add(1)
.ok_or(ListPartitionsError::EntriesExceeded {
limit: max_entries,
requested: usize::MAX,
})?;
if requested_entries > max_entries {
return Err(ListPartitionsError::EntriesExceeded {
limit: max_entries,
requested: requested_entries,
});
}
let requested_name_bytes = name_bytes.checked_add(logical.len()).ok_or(
ListPartitionsError::NameBytesExceeded {
limit: max_name_bytes,
requested: usize::MAX,
},
)?;
if requested_name_bytes > max_name_bytes {
return Err(ListPartitionsError::NameBytesExceeded {
limit: max_name_bytes,
requested: requested_name_bytes,
});
}
name_bytes = requested_name_bytes;
partitions.push(logical);
}
refuse_orphan_remnants(remnants, &data_dirs)?;
partitions.sort_unstable();
if let Some(duplicate) = partitions.windows(2).find(|pair| pair[0] == pair[1]) {
let physical = partition_name::encode_any_namespace(&duplicate[0]).map_or_else(
|_| duplicate[0].clone(),
partition_name::EncodedPartitionName::into_string,
);
return Err(ListPartitionsError::Corrupt {
entry: format!(
"two directories decode to `{}`, one of them `{physical}_data`",
duplicate[0]
),
});
}
Ok(partitions)
}
#[derive(Debug, thiserror::Error)]
pub enum ListPartitionsError {
#[error("storage directory read failed: {0}")]
Storage(std::io::Error),
#[error("`{entry}` has the shape of a partition directory and no logical name")]
Corrupt {
entry: String,
},
#[error("partition listing requested {requested} names, above the limit {limit}")]
EntriesExceeded {
limit: usize,
requested: usize,
},
#[error("partition listing requested {requested} name bytes, above the limit {limit}")]
NameBytesExceeded {
limit: usize,
requested: usize,
},
}
fn refuse_orphan_remnants(
remnants: Vec<String>,
data_dirs: &std::collections::BTreeSet<String>,
) -> Result<(), ListPartitionsError> {
for remnant in remnants {
let Some(base) = remnant_base(&remnant) else {
continue;
};
if data_dirs.contains(base)
|| base.starts_with(REPAIR_STAGE_PREFIX)
|| base.starts_with(REWRITE_STAGE_PREFIX)
|| base.starts_with(EXACT_REWRITE_STAGE_PREFIX)
{
continue;
}
if remnant.ends_with(CHECKPOINT_STORAGE_SUFFIX) {
continue;
}
return Err(ListPartitionsError::Corrupt { entry: remnant });
}
Ok(())
}
fn partition_shaped_remnant(name: &std::ffi::OsStr) -> Option<String> {
let name = name.to_str()?;
OFFSETS_SUFFIXES
.iter()
.chain(std::iter::once(&CHECKPOINT_STORAGE_SUFFIX))
.any(|suffix| name.ends_with(suffix))
.then(|| name.to_owned())
}
fn remnant_base(remnant: &str) -> Option<&str> {
OFFSETS_SUFFIXES
.iter()
.chain(std::iter::once(&CHECKPOINT_STORAGE_SUFFIX))
.find_map(|suffix| remnant.strip_suffix(suffix))
}
async fn partition_presence_in<E>(
context: &E,
storage_dir: &std::path::Path,
encoded: &str,
) -> Result<PartitionPresence, ListPartitionsError>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
if partition_data_dir_exists(storage_dir, encoded)? {
return Ok(PartitionPresence::Present);
}
for suffix in OFFSETS_SUFFIXES {
if directory_present(storage_dir, &format!("{encoded}{suffix}"))? {
return Ok(PartitionPresence::Damaged {
remnant: format!("{encoded}{suffix}"),
});
}
}
let checkpoint = format!("{encoded}{CHECKPOINT_STORAGE_SUFFIX}");
if directory_present(storage_dir, &checkpoint)? {
if recorded_event_count(context, encoded).await? > 0 {
return Ok(PartitionPresence::Damaged {
remnant: checkpoint,
});
}
}
Ok(PartitionPresence::Absent)
}
fn directory_present(
storage_dir: &std::path::Path,
name: &str,
) -> Result<bool, ListPartitionsError> {
match std::fs::symlink_metadata(storage_dir.join(name)) {
Ok(metadata) if metadata.is_dir() => Ok(true),
Ok(_) => Err(ListPartitionsError::Corrupt {
entry: name.to_owned(),
}),
Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(false),
Err(e) => Err(ListPartitionsError::Storage(e)),
}
}
async fn recorded_event_count<E>(context: &E, encoded: &str) -> Result<u64, ListPartitionsError>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
polyc_eventlog::EventCountCheckpoint::open(context.child(presence_label(encoded)), encoded)
.await
.map(|checkpoint| checkpoint.expected_count())
.map_err(|error| ListPartitionsError::Storage(std::io::Error::other(error.to_string())))
}
fn presence_label(encoded: &str) -> &'static str {
static LABELS: std::sync::OnceLock<std::sync::Mutex<HashMap<String, &'static str>>> =
std::sync::OnceLock::new();
let mut labels = LABELS
.get_or_init(|| std::sync::Mutex::new(HashMap::new()))
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if let Some(label) = labels.get(encoded) {
return label;
}
let label: &'static str =
Box::leak(format!("{}_presence", metric_label_body(encoded)).into_boxed_str());
labels.insert(encoded.to_owned(), label);
label
}
fn partition_data_dir_exists(
storage_dir: &std::path::Path,
encoded: &str,
) -> Result<bool, ListPartitionsError> {
directory_present(storage_dir, &format!("{encoded}_data"))
}
fn last_modified_ms_in(storage_dir: &std::path::Path, partition: &PartitionRef) -> Option<u64> {
let data_dir = storage_dir.join(format!("{}_data", partition.storage_key()));
let mut newest: Option<std::time::SystemTime> = None;
for entry in std::fs::read_dir(&data_dir).ok()?.flatten() {
let Ok(modified) = entry.metadata().and_then(|m| m.modified()) else {
continue;
};
newest = Some(match newest {
Some(current) if current >= modified => current,
_ => modified,
});
}
let duration = newest?
.duration_since(std::time::SystemTime::UNIX_EPOCH)
.ok()?;
u64::try_from(duration.as_millis()).ok()
}
#[allow(clippy::too_many_arguments)]
async fn resume_rewrite<E>(
context: &E,
storage_dir: &std::path::Path,
logs: &mut OpenLogs<E>,
checkpoints: &mut HashMap<String, CheckpointFrame>,
count_checkpoints: &mut CountCheckpoints<E>,
mmr_logs: &mut HashMap<String, polyc_mmr::VerifiableLog>,
partition: &str,
stage: &str,
) -> Result<Option<usize>, EventLogError>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
if !logs.contains(stage) && !partition_data_dir_exists(storage_dir, stage).unwrap_or(true) {
return Ok(None);
}
let staged = {
let log = logs.get_or_open(context, stage).await?;
let held = log.replay().await?;
staged_rewrite(&held)
};
let Some((restored, dropped)) = staged else {
return Ok(None);
};
restore_rewrite(
context,
logs,
checkpoints,
count_checkpoints,
mmr_logs,
partition,
&restored,
)
.await?;
destroy_one(
context,
logs,
checkpoints,
count_checkpoints,
mmr_logs,
stage,
DestroyCheckpoint::Clear,
)
.await?;
tracing::info!(
partition = %partition,
dropped,
"partition rewrite resumed from its durable stage (#2323)"
);
Ok(Some(usize::try_from(dropped).unwrap_or(usize::MAX)))
}
#[allow(clippy::too_many_arguments)]
async fn restore_rewrite<E>(
context: &E,
logs: &mut OpenLogs<E>,
checkpoints: &mut HashMap<String, CheckpointFrame>,
count_checkpoints: &mut CountCheckpoints<E>,
mmr_logs: &mut HashMap<String, polyc_mmr::VerifiableLog>,
partition: &str,
restored: &[Event],
) -> Result<u64, EventLogError>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
destroy_one(
context,
logs,
checkpoints,
count_checkpoints,
mmr_logs,
partition,
DestroyCheckpoint::Retain,
)
.await?;
let new_len = if restored.is_empty() {
0
} else {
append_repair_copy(context, logs, partition, restored).await?;
let log = logs.get_or_open(context, partition).await?;
log.len().await
};
match maybe_inject_count_record_failure(partition) {
Some(error) => return Err(error),
None => {
count_checkpoints
.record(context, partition, new_len)
.await?;
}
}
Ok(new_len)
}
#[allow(clippy::too_many_arguments)]
async fn rewrite_one<E>(
context: &E,
storage_dir: &std::path::Path,
logs: &mut OpenLogs<E>,
checkpoints: &mut HashMap<String, CheckpointFrame>,
count_checkpoints: &mut CountCheckpoints<E>,
mmr_logs: &mut HashMap<String, polyc_mmr::VerifiableLog>,
signer: &JournalAttestationSigner,
command_id: &str,
partition: &str,
decide: &(dyn Fn(&Event) -> RewriteDecision + Send + Sync),
) -> Result<usize, EventLogError>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
let stage = rewrite_stage_partition(partition, command_id);
if let Some(dropped) = resume_rewrite(
context,
storage_dir,
logs,
checkpoints,
count_checkpoints,
mmr_logs,
partition,
&stage,
)
.await?
{
return Ok(dropped);
}
let all = {
let log = logs.get_or_open(context, partition).await?;
log.replay().await?
};
let mut survivors: Vec<Event> = Vec::with_capacity(all.len());
let mut dropped = 0usize;
for event in &all {
if event.kind == polyc_eventlog::MMR_SIGNED_ROOT_KIND {
continue;
}
match decide(event) {
RewriteDecision::Keep => survivors.push(event.clone()),
RewriteDecision::Drop => dropped += 1,
RewriteDecision::Replace(payload) => {
survivors.push(Event::with_trust(event.kind.clone(), payload, event.trust));
}
}
}
if dropped == 0
&& survivors
.iter()
.zip(
all.iter()
.filter(|e| e.kind != polyc_eventlog::MMR_SIGNED_ROOT_KIND),
)
.all(|(after, before)| after == before)
{
return Ok(0);
}
let restored = if survivors.is_empty() {
Vec::new()
} else {
let tree = polyc_eventlog::rebuild_from_events(&survivors)?;
let marker = polyc_eventlog::extend_and_sign(&tree, &[], signer)?;
let mut sequence = survivors;
sequence.push(marker);
sequence
};
destroy_one(
context,
logs,
checkpoints,
count_checkpoints,
mmr_logs,
&stage,
DestroyCheckpoint::Clear,
)
.await?;
let mut staged = restored.clone();
staged.push(rewrite_stage_complete(
restored.len() as u64,
dropped as u64,
));
append_repair_copy(context, logs, &stage, &staged).await?;
restore_rewrite(
context,
logs,
checkpoints,
count_checkpoints,
mmr_logs,
partition,
&restored,
)
.await?;
destroy_one(
context,
logs,
checkpoints,
count_checkpoints,
mmr_logs,
&stage,
DestroyCheckpoint::Clear,
)
.await?;
tracing::info!(partition = %partition, dropped, "partition rewritten (erasure)");
Ok(dropped)
}
fn exact_rewrite_content(
all: &[Event],
signer: &JournalAttestationSigner,
decide: &(dyn Fn(&Event) -> RewriteDecision + Send + Sync),
) -> Result<(Vec<Event>, usize), EventLogError> {
let mut survivors = Vec::with_capacity(all.len());
let mut dropped = 0_usize;
for event in all {
if event.kind == polyc_eventlog::MMR_SIGNED_ROOT_KIND {
continue;
}
match decide(event) {
RewriteDecision::Keep => survivors.push(event.clone()),
RewriteDecision::Drop => dropped = dropped.saturating_add(1),
RewriteDecision::Replace(payload) => {
survivors.push(Event::with_trust(event.kind.clone(), payload, event.trust));
}
}
}
if survivors.is_empty() {
return Ok((Vec::new(), dropped));
}
let tree = polyc_eventlog::rebuild_from_events(&survivors)?;
let root = polyc_eventlog::extend_and_sign(&tree, &[], signer)?;
survivors.push(root);
Ok((survivors, dropped))
}
fn inspect_exact_rewrite(
partition: &str,
command_id: &str,
all: &[Event],
signer: &JournalAttestationSigner,
decide: &(dyn Fn(&Event) -> RewriteDecision + Send + Sync),
) -> Result<(RewritePlan, Vec<Event>, usize), EventLogError> {
let before = physical_fingerprint(all);
let (restored, dropped) = exact_rewrite_content(all, signer, decide)?;
let after = physical_fingerprint(&restored);
let stage_binding = exact_rewrite_binding(partition, command_id, &before, &after);
Ok((
RewritePlan {
before,
after,
stage_binding,
},
restored,
dropped,
))
}
#[allow(clippy::too_many_arguments)]
async fn apply_exact_rewrite<E>(
context: &E,
logs: &mut OpenLogs<E>,
checkpoints: &mut HashMap<String, CheckpointFrame>,
count_checkpoints: &mut CountCheckpoints<E>,
mmr_logs: &mut HashMap<String, polyc_mmr::VerifiableLog>,
signer: &JournalAttestationSigner,
partition: &str,
command_id: &str,
expected: RewritePlan,
decide: &(dyn Fn(&Event) -> RewriteDecision + Send + Sync),
) -> Result<RewriteApplication, EventLogError>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
let RewritePlan {
before,
after,
stage_binding,
} = expected;
if exact_rewrite_binding(partition, command_id, &before, &after) != stage_binding {
return Ok(RewriteApplication::Changed);
}
let stage = exact_rewrite_stage_partition(partition, &stage_binding);
let current = {
let log = logs.get_or_open(context, partition).await?;
log.replay().await?
};
let current_fingerprint = physical_fingerprint(¤t);
if current_fingerprint == after {
destroy_one(
context,
logs,
checkpoints,
count_checkpoints,
mmr_logs,
&stage,
DestroyCheckpoint::Clear,
)
.await?;
return Ok(RewriteApplication::AlreadyApplied);
}
let (restored, dropped) = if current_fingerprint == before {
let (observed, restored, dropped) =
inspect_exact_rewrite(partition, command_id, ¤t, signer, decide)?;
if observed
!= (RewritePlan {
before,
after,
stage_binding,
})
{
return Ok(RewriteApplication::Changed);
}
destroy_one(
context,
logs,
checkpoints,
count_checkpoints,
mmr_logs,
&stage,
DestroyCheckpoint::Clear,
)
.await?;
let mut staged = restored.clone();
staged.push(exact_rewrite_stage_complete(
before,
after,
stage_binding,
restored.len() as u64,
dropped as u64,
));
append_repair_copy(context, logs, &stage, &staged).await?;
(restored, dropped)
} else {
let held = {
let log = logs.get_or_open(context, &stage).await?;
log.replay().await?
};
let Some((restored, dropped)) =
exact_staged_rewrite(&held, &before, &after, &stage_binding)
else {
return Ok(RewriteApplication::Changed);
};
(restored, usize::try_from(dropped).unwrap_or(usize::MAX))
};
restore_rewrite(
context,
logs,
checkpoints,
count_checkpoints,
mmr_logs,
partition,
&restored,
)
.await?;
destroy_one(
context,
logs,
checkpoints,
count_checkpoints,
mmr_logs,
&stage,
DestroyCheckpoint::Clear,
)
.await?;
Ok(RewriteApplication::Applied { dropped })
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum DestroyCheckpoint {
Clear,
Retain,
}
#[allow(clippy::too_many_arguments)]
async fn destroy_one<E>(
context: &E,
logs: &mut OpenLogs<E>,
checkpoints: &mut HashMap<String, CheckpointFrame>,
count_checkpoints: &mut CountCheckpoints<E>,
mmr_logs: &mut HashMap<String, polyc_mmr::VerifiableLog>,
partition: &str,
checkpoint: DestroyCheckpoint,
) -> Result<(), EventLogError>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
if checkpoint == DestroyCheckpoint::Clear {
count_checkpoints.clear(context, partition).await?;
}
logs.get_or_open(context, partition).await?;
let label = logs.log_label(partition);
let log = logs.take(partition).expect("just opened via get_or_open");
if let Some(err) = maybe_inject_destroy_failure(partition) {
return Err(err);
}
{
let emptied = Box::pin(log.reset(context.child(label))).await?;
Box::pin(emptied.reclaim()).await?;
}
checkpoints.remove(partition);
mmr_logs.remove(partition);
count_checkpoints.forget(partition);
tracing::info!(partition = %partition, "partition erased (reset, then reclaimed)");
Ok(())
}
#[cfg(test)]
fn maybe_inject_append_failure(partition: &str, event_index: usize) -> Option<EventLogError> {
tests::take_armed_append_failure(partition, event_index)
.or_else(|| take_armed_append_failure_at(partition, event_index))
}
#[cfg(all(not(test), feature = "test-util"))]
fn maybe_inject_append_failure(partition: &str, event_index: usize) -> Option<EventLogError> {
take_armed_append_failure_at(partition, event_index)
}
#[cfg(not(any(test, feature = "test-util")))]
const fn maybe_inject_append_failure(
_partition: &str,
_event_index: usize,
) -> Option<EventLogError> {
None
}
#[cfg(test)]
fn maybe_inject_commit_failure(partition: &str) -> Option<EventLogError> {
tests::take_armed_commit_failure(partition)
}
#[cfg(any(test, feature = "test-util"))]
static ARMED_APPEND_FAILURES_AT: std::sync::OnceLock<
std::sync::Mutex<std::collections::HashMap<String, usize>>,
> = std::sync::OnceLock::new();
#[cfg(feature = "test-util")]
fn arm_append_failure_at(partition: &str, event_index: usize) {
ARMED_APPEND_FAILURES_AT
.get_or_init(|| std::sync::Mutex::new(std::collections::HashMap::new()))
.lock()
.expect("lock poisoned")
.insert(partition.to_owned(), event_index);
}
#[cfg(any(test, feature = "test-util"))]
fn take_armed_append_failure_at(partition: &str, event_index: usize) -> Option<EventLogError> {
let lock = ARMED_APPEND_FAILURES_AT.get()?;
let matched = {
let mut armed = lock.lock().expect("lock poisoned");
if armed.get(partition) == Some(&event_index) {
armed.remove(partition);
true
} else {
false
}
};
matched.then(|| {
EventLogError::Journal(commonware_storage::journal::Error::Corruption(
"injected test failure (#1712 fault-injection seam)".to_owned(),
))
})
}
#[cfg(any(test, feature = "test-util"))]
static ARMED_ROOT_APPEND_FAILURES: std::sync::OnceLock<
std::sync::Mutex<std::collections::HashMap<String, usize>>,
> = std::sync::OnceLock::new();
#[cfg(any(test, feature = "test-util"))]
fn arm_root_append_failure(partition: &str) {
ARMED_ROOT_APPEND_FAILURES
.get_or_init(|| std::sync::Mutex::new(std::collections::HashMap::new()))
.lock()
.expect("lock poisoned")
.entry(partition.to_owned())
.and_modify(|count| *count = count.saturating_add(1))
.or_insert(1);
}
#[cfg(any(test, feature = "test-util"))]
fn take_armed_root_append_failure(partition: &str) -> Option<EventLogError> {
let lock = ARMED_ROOT_APPEND_FAILURES.get()?;
let matched = {
let mut armed = lock.lock().expect("lock poisoned");
match armed.get_mut(partition) {
Some(count) if *count > 1 => {
*count -= 1;
true
}
Some(_) => {
armed.remove(partition);
true
}
None => false,
}
};
matched.then(|| {
EventLogError::Journal(commonware_storage::journal::Error::Corruption(
"injected root-append failure".to_owned(),
))
})
}
#[cfg(any(test, feature = "test-util"))]
fn maybe_inject_root_append_failure(partition: &str) -> Option<EventLogError> {
take_armed_root_append_failure(partition)
}
#[cfg(test)]
fn maybe_inject_mmr_preparation_failure(partition: &str) -> Option<EventLogError> {
tests::take_armed_mmr_preparation_failure(partition)
}
#[cfg(not(any(test, feature = "test-util")))]
const fn maybe_inject_root_append_failure(_partition: &str) -> Option<EventLogError> {
None
}
#[cfg(not(test))]
const fn maybe_inject_mmr_preparation_failure(_partition: &str) -> Option<EventLogError> {
None
}
#[cfg(not(test))]
const fn maybe_inject_commit_failure(_partition: &str) -> Option<EventLogError> {
None
}
#[cfg(test)]
fn maybe_inject_count_record_failure(partition: &str) -> Option<EventLogError> {
tests::take_armed_count_record_failure(partition)
}
#[cfg(not(test))]
const fn maybe_inject_count_record_failure(_partition: &str) -> Option<EventLogError> {
None
}
#[cfg(test)]
fn maybe_inject_destroy_failure(partition: &str) -> Option<EventLogError> {
tests::take_armed_destroy_failure(partition)
}
#[cfg(not(test))]
const fn maybe_inject_destroy_failure(_partition: &str) -> Option<EventLogError> {
None
}
#[allow(clippy::too_many_arguments)]
async fn append_batch_one<E>(
context: &E,
logs: &mut OpenLogs<E>,
checkpoints: &mut HashMap<String, CheckpointFrame>,
count_checkpoints: &mut CountCheckpoints<E>,
mmr_logs: &mut HashMap<String, polyc_mmr::VerifiableLog>,
signer: &JournalAttestationSigner,
partition: &str,
events: &[Event],
) -> Result<Vec<u64>, AppendError>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
let log = logs
.get_or_open(context, partition)
.await
.map_err(AppendError::Log)?;
let prepared_root = if events.is_empty() {
None
} else {
match prepare_signed_root(log, mmr_logs, signer, partition, events).await {
Ok(marker) => Some(marker),
Err(error) => {
metrics::record_attestation(error.outcome());
logs.take(partition);
mmr_logs.remove(partition);
return Err(AppendError::Log(error.into_eventlog_error()));
}
}
};
let mut positions = Vec::with_capacity(events.len());
for (idx, event) in events.iter().enumerate() {
let append_result = match maybe_inject_append_failure(partition, idx) {
Some(err) => Err(err),
None => log.append(event).await,
};
match append_result {
Ok(pos) => positions.push(pos),
Err(err) => {
match log.commit().await {
Ok(()) if !positions.is_empty() => {
metrics::record_attestation(
metrics::AttestationOutcome::UnattestedPartialBatch,
);
}
Ok(()) => {}
Err(commit_err) => {
tracing::warn!(error = %commit_err, "commit of partial batch failed");
}
}
logs.take(partition);
mmr_logs.remove(partition);
return Err(AppendError::AppendOutcomeUnknown(err));
}
}
}
let attestation = prepared_root
.is_some()
.then_some(metrics::AttestationOutcome::Signed);
if let Some(marker) = prepared_root
&& let Err(error) = append_prepared_root(log, partition, &marker).await
{
logs.take(partition);
mmr_logs.remove(partition);
return Err(AppendError::AttestationOutcomeUnknown(error));
}
let commit_result = match maybe_inject_commit_failure(partition) {
Some(err) => Err(err),
None => log.commit().await,
};
if let Err(err) = commit_result {
logs.take(partition);
mmr_logs.remove(partition);
return Err(AppendError::AttestationOutcomeUnknown(err));
}
if let Some(outcome) = attestation {
metrics::record_attestation(outcome);
}
if !events.is_empty() {
record_event_count_floor(context, count_checkpoints, partition, log.len().await).await;
}
if let Some((pos, event)) = events
.iter()
.zip(&positions)
.rev()
.find(|(event, _)| parse_kind(&event.kind).0 == BASE_COMPACTION_CHECKPOINT)
.map(|(event, pos)| (*pos, event))
{
checkpoints.insert(
partition.to_owned(),
CheckpointFrame {
tail_offset: pos + 1,
payload: event.payload.clone(),
},
);
}
Ok(positions)
}
async fn record_event_count_floor<E>(
context: &E,
count_checkpoints: &mut CountCheckpoints<E>,
partition: &str,
new_len: u64,
) where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
match count_checkpoints.record(context, partition, new_len).await {
Ok(()) => {}
Err(err) => {
tracing::warn!(
partition = %partition,
error = %err,
"event-count checkpoint record failed (#799); truncation detection stays at its last recorded value until the next successful commit"
);
}
}
}
#[derive(Debug, thiserror::Error)]
enum MmrPreparationError {
#[error("mmr rebuild: {0}")]
Rebuild(EventLogError),
#[error("mmr extend/sign: {0}")]
ExtendSign(polyc_eventlog::IntegrityError),
}
impl MmrPreparationError {
fn into_eventlog_error(self) -> EventLogError {
match self {
Self::Rebuild(error) => error,
Self::ExtendSign(error) => EventLogError::Integrity(error),
}
}
const fn outcome(&self) -> metrics::AttestationOutcome {
match self {
Self::Rebuild(_) => metrics::AttestationOutcome::RebuildFailed,
Self::ExtendSign(_) => metrics::AttestationOutcome::ExtendSignFailed,
}
}
}
async fn append_prepared_root<E>(
log: &polyc_eventlog::EventLog<E>,
partition: &str,
marker: &Event,
) -> Result<(), EventLogError>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
let append_result = match maybe_inject_root_append_failure(partition) {
Some(error) => Err(error),
None => log.append(marker).await.map(|_| ()),
};
let Err(error) = append_result else {
return Ok(());
};
if let Err(commit_error) = log.commit().await {
tracing::warn!(
error = %commit_error,
"commit after uncertain root-marker append failed"
);
}
metrics::record_attestation(metrics::AttestationOutcome::MarkerAppendFailed);
Err(error)
}
async fn prepare_signed_root<E>(
log: &EventLog<E>,
mmr_logs: &mut HashMap<String, polyc_mmr::VerifiableLog>,
signer: &JournalAttestationSigner,
partition: &str,
events: &[Event],
) -> Result<Event, MmrPreparationError>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler,
{
if let Some(error) = maybe_inject_mmr_preparation_failure(partition) {
return Err(MmrPreparationError::Rebuild(error));
}
if !mmr_logs.contains_key(partition) {
let all = log.replay().await.map_err(MmrPreparationError::Rebuild)?;
let rebuilt =
polyc_eventlog::rebuild_from_events(&all).map_err(MmrPreparationError::ExtendSign)?;
mmr_logs.insert(partition.to_owned(), rebuilt);
}
let mmr_log = mmr_logs.get(partition).expect("just ensured present above");
polyc_eventlog::extend_and_sign(mmr_log, events, signer)
.map_err(MmrPreparationError::ExtendSign)
}
enum EnsureFailure {
Definite(EventLogError),
Ambiguous(EventLogError),
}
#[allow(clippy::too_many_arguments)]
async fn ensure_partition_attested_one<E>(
context: &E,
storage_dir: &std::path::Path,
logs: &mut OpenLogs<E>,
count_checkpoints: &mut CountCheckpoints<E>,
mmr_logs: &mut HashMap<String, polyc_mmr::VerifiableLog>,
signer: &JournalAttestationSigner,
trust: &polyc_crypto::signing_role::RoleTrustSet<
polyc_crypto::signing_role::JournalAttestationRole,
>,
partition: &str,
) -> Result<Option<Event>, AppendError>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
if !logs.contains(partition)
&& !partition_data_dir_exists(storage_dir, partition).map_err(AppendError::Listing)?
{
return Ok(None);
}
let result = async {
let log = logs
.get_or_open(context, partition)
.await
.map_err(EnsureFailure::Definite)?;
let events = log.replay().await.map_err(EnsureFailure::Definite)?;
polyc_eventlog::verify_replay_with_trust(&events, trust)
.map_err(EventLogError::from)
.map_err(EnsureFailure::Definite)?;
let rebuilt = polyc_eventlog::rebuild_from_events(&events)
.map_err(EventLogError::from)
.map_err(EnsureFailure::Definite)?;
let leaf_count = rebuilt
.leaf_count()
.map_err(polyc_eventlog::IntegrityError::from)
.map_err(EventLogError::from)
.map_err(EnsureFailure::Definite)?;
if leaf_count == 0 {
return Ok((None, rebuilt, log.len().await));
}
let exact = events.iter().rev().find_map(|event| {
if event.kind != polyc_eventlog::MMR_SIGNED_ROOT_KIND {
return None;
}
serde_json::from_slice::<polyc_mmr::SignedRoot>(&event.payload)
.ok()
.filter(|root| root.leaf_count == leaf_count)
.map(|_| event.clone())
});
let marker = if let Some(marker) = exact {
marker
} else {
let marker = polyc_eventlog::extend_and_sign(&rebuilt, &[], signer)
.map_err(EventLogError::from)
.map_err(EnsureFailure::Definite)?;
match maybe_inject_root_append_failure(partition) {
Some(error) => return Err(EnsureFailure::Ambiguous(error)),
None => log
.append(&marker)
.await
.map_err(EnsureFailure::Ambiguous)?,
};
match maybe_inject_commit_failure(partition) {
Some(error) => return Err(EnsureFailure::Ambiguous(error)),
None => log.commit().await.map_err(EnsureFailure::Ambiguous)?,
}
marker
};
let new_len = log.len().await;
Ok::<_, EnsureFailure>((Some(marker), rebuilt, new_len))
}
.await;
let (marker, rebuilt, new_len) = match result {
Ok(value) => value,
Err(EnsureFailure::Definite(error)) => {
logs.take(partition);
mmr_logs.remove(partition);
return Err(AppendError::Log(error));
}
Err(EnsureFailure::Ambiguous(error)) => {
logs.take(partition);
mmr_logs.remove(partition);
return Err(AppendError::AttestationOutcomeUnknown(error));
}
};
mmr_logs.insert(partition.to_owned(), rebuilt);
if let Err(error) = count_checkpoints.record(context, partition, new_len).await {
tracing::warn!(
partition = %partition,
error = %error,
"event-count checkpoint record failed after attestation completion"
);
}
Ok(marker)
}
async fn replay_verified_one<E>(
context: &E,
logs: &mut OpenLogs<E>,
count_checkpoints: &mut CountCheckpoints<E>,
journal_attestation_trust: &polyc_crypto::signing_role::RoleTrustSet<
polyc_crypto::signing_role::JournalAttestationRole,
>,
partition: &str,
) -> Result<Vec<Event>, VerifyError>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
let log = logs.get_or_open(context, partition).await?;
let events = log.replay().await?;
let expected = count_checkpoints.expected_count(context, partition).await?;
let actual = events.len() as u64;
if actual < expected {
return Err(VerifyError::TruncatedReplay { expected, actual });
}
polyc_eventlog::verify_replay_with_trust(&events, journal_attestation_trust)?;
Ok(events)
}
#[derive(Debug)]
struct RepairEffect {
quarantined: Vec<QuarantinedItem>,
rewrote: bool,
}
#[derive(Clone)]
struct RepairScan {
readable: Vec<(u64, Event)>,
quarantined: Vec<QuarantinedItem>,
before: [u8; 32],
after: [u8; 32],
standing: polyc_eventlog::RootStanding,
}
pub(crate) const REPAIR_STAGE_PREFIX: &str = "state-repair-stage-";
pub(crate) const REWRITE_STAGE_PREFIX: &str = "state-rewrite-stage-";
const EXACT_REWRITE_STAGE_PREFIX: &str = "state-exact-rewrite-stage-";
const EXACT_REWRITE_STAGE_COMPLETE_KIND: &str = "__polyc_state_exact_rewrite_stage_complete";
const REWRITE_STAGE_COMPLETE_KIND: &str = "__polyc_state_rewrite_stage_complete";
fn rewrite_stage_partition(partition: &str, command_id: &str) -> String {
use sha2::Digest as _;
let mut hash = sha2::Sha256::new();
hash.update(b"polychrome/eventlog-rewrite-stage/v1");
hash.update((partition.len() as u64).to_be_bytes());
hash.update(partition.as_bytes());
hash.update((command_id.len() as u64).to_be_bytes());
hash.update(command_id.as_bytes());
format!("{REWRITE_STAGE_PREFIX}{}", hex::encode(hash.finalize()))
}
fn rewrite_stage_complete(count: u64, dropped: u64) -> Event {
let mut payload = Vec::with_capacity(16);
payload.extend_from_slice(&count.to_be_bytes());
payload.extend_from_slice(&dropped.to_be_bytes());
Event::new(REWRITE_STAGE_COMPLETE_KIND, payload)
}
fn staged_rewrite(events: &[Event]) -> Option<(Vec<Event>, u64)> {
let (marker, survivors) = events.split_last()?;
if marker.kind != REWRITE_STAGE_COMPLETE_KIND || marker.payload.len() != 16 {
return None;
}
let count = u64::from_be_bytes(marker.payload[..8].try_into().ok()?);
let dropped = u64::from_be_bytes(marker.payload[8..].try_into().ok()?);
(count == survivors.len() as u64).then(|| (survivors.to_vec(), dropped))
}
fn physical_fingerprint(events: &[Event]) -> [u8; 32] {
let indexed = events
.iter()
.cloned()
.enumerate()
.map(|(position, event)| (position as u64, event))
.collect::<Vec<_>>();
repair_fingerprint(&indexed, &[], events.len() as u64)
}
fn exact_rewrite_binding(
partition: &str,
command_id: &str,
before: &[u8; 32],
after: &[u8; 32],
) -> [u8; 32] {
use sha2::Digest as _;
let mut hash = sha2::Sha256::new();
hash.update(b"polychrome/exact-rewrite-stage/v1");
hash.update((partition.len() as u64).to_be_bytes());
hash.update(partition.as_bytes());
hash.update((command_id.len() as u64).to_be_bytes());
hash.update(command_id.as_bytes());
hash.update(before);
hash.update(after);
hash.finalize().into()
}
fn exact_rewrite_stage_partition(partition: &str, binding: &[u8; 32]) -> String {
use sha2::Digest as _;
let mut hash = sha2::Sha256::new();
hash.update(b"polychrome/exact-rewrite-stage-name/v1");
hash.update((partition.len() as u64).to_be_bytes());
hash.update(partition.as_bytes());
hash.update(binding);
format!(
"{EXACT_REWRITE_STAGE_PREFIX}{}",
hex::encode(hash.finalize())
)
}
fn exact_rewrite_stage_complete(
before: [u8; 32],
after: [u8; 32],
binding: [u8; 32],
count: u64,
dropped: u64,
) -> Event {
let mut payload = Vec::with_capacity(112);
payload.extend_from_slice(&before);
payload.extend_from_slice(&after);
payload.extend_from_slice(&binding);
payload.extend_from_slice(&count.to_be_bytes());
payload.extend_from_slice(&dropped.to_be_bytes());
Event::new(EXACT_REWRITE_STAGE_COMPLETE_KIND, payload)
}
fn exact_staged_rewrite(
events: &[Event],
before: &[u8; 32],
after: &[u8; 32],
binding: &[u8; 32],
) -> Option<(Vec<Event>, u64)> {
let (marker, survivors) = events.split_last()?;
if marker.kind != EXACT_REWRITE_STAGE_COMPLETE_KIND || marker.payload.len() != 112 {
return None;
}
if marker.payload[..32] != before[..]
|| marker.payload[32..64] != after[..]
|| marker.payload[64..96] != binding[..]
|| marker.payload[96..104] != (survivors.len() as u64).to_be_bytes()
|| physical_fingerprint(survivors) != *after
{
return None;
}
let dropped = u64::from_be_bytes(marker.payload[104..112].try_into().ok()?);
Some((survivors.to_vec(), dropped))
}
const REPAIR_STAGE_COMPLETE_KIND: &str = "__polyc_state_repair_stage_complete";
fn repair_stage_partition(partition: &str, before: &[u8; 32]) -> String {
use sha2::Digest as _;
let mut hash = sha2::Sha256::new();
hash.update(b"polychrome/eventlog-repair-stage/v1");
hash.update((partition.len() as u64).to_be_bytes());
hash.update(partition.as_bytes());
hash.update(before);
format!("{REPAIR_STAGE_PREFIX}{}", hex::encode(hash.finalize()))
}
fn repair_stage_complete(before: [u8; 32], after: [u8; 32], count: u64) -> Event {
let mut payload = Vec::with_capacity(72);
payload.extend_from_slice(&before);
payload.extend_from_slice(&after);
payload.extend_from_slice(&count.to_be_bytes());
Event::new(REPAIR_STAGE_COMPLETE_KIND, payload)
}
async fn repair_window_is_open<E>(
context: &E,
count_checkpoints: &mut CountCheckpoints<E>,
partition: &str,
) -> Result<bool, EventLogError>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
Ok(count_checkpoints.expected_count(context, partition).await? == 0)
}
fn staged_survivors(events: &[Event], before: &[u8; 32], after: &[u8; 32]) -> Option<Vec<Event>> {
let (marker, survivors) = events.split_last()?;
if marker.kind != REPAIR_STAGE_COMPLETE_KIND || marker.payload.len() != 72 {
return None;
}
if marker.payload[..32] != before[..]
|| marker.payload[32..64] != after[..]
|| marker.payload[64..] != (survivors.len() as u64).to_be_bytes()
{
return None;
}
let indexed = survivors
.iter()
.enumerate()
.map(|(position, event)| (position as u64, event.clone()))
.collect::<Vec<_>>();
(repair_fingerprint(&indexed, &[], survivors.len() as u64) == *after)
.then(|| survivors.to_vec())
}
impl RepairScan {
const fn will_rewrite(&self) -> bool {
repair_rewrites(self.quarantined.is_empty(), self.standing)
}
fn plan(&self) -> RepairPlan {
RepairPlan {
before: self.before,
after: self.after,
will_rewrite: self.will_rewrite(),
readable: self
.readable
.iter()
.map(|(_, event)| event.clone())
.collect(),
quarantined: self.quarantined.clone(),
}
}
}
const fn repair_rewrites(
quarantine_is_empty: bool,
standing: polyc_eventlog::RootStanding,
) -> bool {
!quarantine_is_empty || !matches!(standing, polyc_eventlog::RootStanding::Holds)
}
fn repair_fingerprint(
readable: &[(u64, Event)],
quarantined: &[QuarantinedItem],
expected: u64,
) -> [u8; 32] {
use sha2::Digest as _;
let mut hash = sha2::Sha256::new();
hash.update(b"polychrome/eventlog-repair-state/v1");
hash.update(expected.to_be_bytes());
hash.update((readable.len() as u64).to_be_bytes());
for (position, event) in readable {
hash.update(position.to_be_bytes());
hash.update((event.kind.len() as u64).to_be_bytes());
hash.update(event.kind.as_bytes());
hash.update([event.trust.as_u8()]);
hash.update((event.payload.len() as u64).to_be_bytes());
hash.update(&event.payload);
}
hash.update((quarantined.len() as u64).to_be_bytes());
for item in quarantined {
hash.update(item.position.to_be_bytes());
}
hash.finalize().into()
}
fn repaired_content(
readable: &[Event],
signer: &JournalAttestationSigner,
replacement: Option<&RepairEventReplacement>,
) -> Result<(Vec<Event>, Option<polyc_mmr::VerifiableLog>), EventLogError> {
let mut survivors: Vec<Event> = readable
.iter()
.filter(|event| event.kind != polyc_eventlog::MMR_SIGNED_ROOT_KIND)
.cloned()
.collect();
if let Some(replacement) = replacement {
let mut replaced = 0_u8;
for event in &mut survivors {
if event.kind == replacement.kind {
event.payload.clone_from(&replacement.payload);
replaced = replaced.saturating_add(1);
}
}
if replaced != 1 {
return Err(EventLogError::RepairReplacement(format!(
"expected one {} event, found {replaced}",
replacement.kind
)));
}
}
if survivors.is_empty() {
return Ok((Vec::new(), None));
}
let tree = polyc_eventlog::rebuild_from_events(&survivors)?;
let marker = polyc_eventlog::extend_and_sign(&tree, &[], signer)?;
let mut content = survivors;
content.push(marker);
Ok((content, Some(tree)))
}
async fn inspect_repair_one<E>(
context: &E,
logs: &mut OpenLogs<E>,
count_checkpoints: &mut CountCheckpoints<E>,
signer: &JournalAttestationSigner,
trust: &polyc_crypto::signing_role::RoleTrustSet<
polyc_crypto::signing_role::JournalAttestationRole,
>,
partition: &str,
replacement: Option<&RepairEventReplacement>,
) -> Result<RepairScan, EventLogError>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
let (ok, mut quarantined) = {
let log = logs.get_or_open(context, partition).await?;
log.replay_quarantining().await?
};
let visible = (ok.len() + quarantined.len()) as u64;
let expected = count_checkpoints.expected_count(context, partition).await?;
if visible < expected {
for position in visible..expected {
quarantined.push(QuarantinedItem {
position,
error: "position lost to the storage layer's silent truncation of the active \
section (torn write or corruption/tampering); original content is \
unrecoverable — see docs/operations/eventlog-disaster-recovery.md"
.to_owned(),
});
}
}
let before = repair_fingerprint(&ok, &quarantined, expected);
let raw: Vec<Event> = ok.iter().map(|(_, event)| event.clone()).collect();
let standing = polyc_eventlog::root_standing_with_trust(&raw, trust)?;
let replacement = repair_rewrites(quarantined.is_empty(), standing)
.then_some(replacement)
.flatten();
let (restored, _) = repaired_content(&raw, signer, replacement)?;
let survivors: Vec<(u64, Event)> = restored
.iter()
.enumerate()
.map(|(position, event)| (position as u64, event.clone()))
.collect();
let survivor_count = survivors.len() as u64;
let after = repair_fingerprint(&survivors, &[], survivor_count);
Ok(RepairScan {
readable: ok,
quarantined,
before,
after,
standing,
})
}
async fn append_repair_copy<E>(
context: &E,
logs: &mut OpenLogs<E>,
partition: &str,
events: &[Event],
) -> Result<(), EventLogError>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
let log = logs.get_or_open(context, partition).await?;
for (index, event) in events.iter().enumerate() {
let append = match maybe_inject_append_failure(partition, index) {
Some(error) => Err(error),
None => log.append(event).await.map(|_| ()),
};
if let Err(error) = append {
let _ = log.commit().await;
logs.take(partition);
return Err(error);
}
}
let commit = match maybe_inject_commit_failure(partition) {
Some(error) => Err(error),
None => log.commit().await,
};
if let Err(error) = commit {
logs.take(partition);
return Err(error);
}
Ok(())
}
#[allow(clippy::too_many_arguments)]
async fn apply_repair_scan<E>(
context: &E,
logs: &mut OpenLogs<E>,
checkpoints: &mut HashMap<String, CheckpointFrame>,
count_checkpoints: &mut CountCheckpoints<E>,
mmr_logs: &mut HashMap<String, polyc_mmr::VerifiableLog>,
signer: &JournalAttestationSigner,
partition: &str,
scan: RepairScan,
replacement: Option<&RepairEventReplacement>,
) -> Result<RepairEffect, EventLogError>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
if scan.standing == polyc_eventlog::RootStanding::Contradicted {
return Err(EventLogError::RepairRefused);
}
if !repair_rewrites(scan.quarantined.is_empty(), scan.standing) {
return Ok(RepairEffect {
quarantined: scan.quarantined,
rewrote: false,
});
}
debug_assert!(
repair_rewrites(scan.quarantined.is_empty(), scan.standing),
"a non-rewriting repair returned above, so this point implies a rewrite"
);
let raw: Vec<Event> = scan.readable.into_iter().map(|(_, event)| event).collect();
let (survivors, tree) = repaired_content(&raw, signer, replacement)?;
let survivor_count = survivors.len() as u64;
let stage = repair_stage_partition(partition, &scan.before);
let staged = {
let log = logs.get_or_open(context, &stage).await?;
let held = log.replay().await?;
staged_survivors(&held, &scan.before, &scan.after)
};
let (survivors, tree) = if let Some(staged) = staged {
(staged, None)
} else {
destroy_one(
context,
logs,
checkpoints,
count_checkpoints,
mmr_logs,
&stage,
DestroyCheckpoint::Clear,
)
.await?;
let mut staged = survivors.clone();
staged.push(repair_stage_complete(
scan.before,
scan.after,
survivor_count,
));
append_repair_copy(context, logs, &stage, &staged).await?;
(survivors, tree)
};
destroy_one(
context,
logs,
checkpoints,
count_checkpoints,
mmr_logs,
partition,
DestroyCheckpoint::Clear,
)
.await?;
let restored_count = survivors.len() as u64;
if !survivors.is_empty() {
append_repair_copy(context, logs, partition, &survivors).await?;
if let Some(tree) = tree {
mmr_logs.insert(partition.to_owned(), tree);
}
}
count_checkpoints
.record(context, partition, restored_count)
.await?;
if let Err(error) = destroy_one(
context,
logs,
checkpoints,
count_checkpoints,
mmr_logs,
&stage,
DestroyCheckpoint::Clear,
)
.await
{
tracing::warn!(
partition = %partition,
%error,
"a completed repair could not remove its survivor stage; the stage is inert now that \
the partition is whole, and the window test keeps it from being replayed over one \
that has moved on"
);
}
tracing::warn!(
partition = %partition,
quarantined = scan.quarantined.len(),
"partition repaired (#799): corrupted events quarantined so the rest replays again"
);
Ok(RepairEffect {
quarantined: scan.quarantined,
rewrote: true,
})
}
async fn destroy_stage_if_present<E>(
context: &E,
storage_dir: &std::path::Path,
logs: &mut OpenLogs<E>,
checkpoints: &mut HashMap<String, CheckpointFrame>,
count_checkpoints: &mut CountCheckpoints<E>,
mmr_logs: &mut HashMap<String, polyc_mmr::VerifiableLog>,
stage: &str,
) -> Result<(), AppendError>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
if !logs.contains(stage) && !partition_data_dir_exists(storage_dir, stage)? {
return Ok(());
}
destroy_one(
context,
logs,
checkpoints,
count_checkpoints,
mmr_logs,
stage,
DestroyCheckpoint::Clear,
)
.await
.map_err(AppendError::Log)
}
#[allow(clippy::too_many_arguments)]
async fn repair_settled_one<E>(
context: &E,
storage_dir: &std::path::Path,
logs: &mut OpenLogs<E>,
count_checkpoints: &mut CountCheckpoints<E>,
signer: &JournalAttestationSigner,
trust: &polyc_crypto::signing_role::RoleTrustSet<
polyc_crypto::signing_role::JournalAttestationRole,
>,
partition: &str,
expected: RepairFingerprints,
) -> Result<bool, AppendError>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
let RepairFingerprints { before, after } = expected;
let stage = repair_stage_partition(partition, &before);
if !logs.contains(&stage) && !partition_data_dir_exists(storage_dir, &stage).unwrap_or(true) {
return Ok(true);
}
let held = {
let log = logs
.get_or_open(context, &stage)
.await
.map_err(AppendError::Log)?;
log.replay().await.map_err(AppendError::Log)?
};
if staged_survivors(&held, &before, &after).is_none() {
return Ok(true);
}
let scan = inspect_repair_one(
context,
logs,
count_checkpoints,
signer,
trust,
partition,
None,
)
.await
.map_err(AppendError::Log)?;
let window_open = repair_window_is_open(context, count_checkpoints, partition)
.await
.map_err(AppendError::Log)?;
Ok(scan.before == after || scan.before == before || !window_open)
}
#[allow(clippy::too_many_arguments)]
async fn sweep_finished_repair_stage<E>(
context: &E,
logs: &mut OpenLogs<E>,
checkpoints: &mut HashMap<String, CheckpointFrame>,
count_checkpoints: &mut CountCheckpoints<E>,
mmr_logs: &mut HashMap<String, polyc_mmr::VerifiableLog>,
partition: &str,
before: &[u8; 32],
) where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
let stage = repair_stage_partition(partition, before);
if let Err(error) = destroy_one(
context,
logs,
checkpoints,
count_checkpoints,
mmr_logs,
&stage,
DestroyCheckpoint::Clear,
)
.await
{
tracing::warn!(
partition = %partition,
%error,
"a finished repair could not remove its survivor stage"
);
}
}
#[derive(Clone, Copy)]
struct RepairFingerprints {
before: [u8; 32],
after: [u8; 32],
}
#[allow(clippy::too_many_arguments)]
async fn apply_repair_conditionally<E>(
context: &E,
logs: &mut OpenLogs<E>,
checkpoints: &mut HashMap<String, CheckpointFrame>,
count_checkpoints: &mut CountCheckpoints<E>,
mmr_logs: &mut HashMap<String, polyc_mmr::VerifiableLog>,
signer: &JournalAttestationSigner,
trust: &polyc_crypto::signing_role::RoleTrustSet<
polyc_crypto::signing_role::JournalAttestationRole,
>,
partition: &str,
expected: RepairFingerprints,
replacement: Option<&RepairEventReplacement>,
) -> Result<(RepairApplication, usize), EventLogError>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
let RepairFingerprints { before, after } = expected;
let scan = inspect_repair_one(
context,
logs,
count_checkpoints,
signer,
trust,
partition,
replacement,
)
.await?;
if scan.before == after && scan.standing == polyc_eventlog::RootStanding::Holds {
sweep_finished_repair_stage(
context,
logs,
checkpoints,
count_checkpoints,
mmr_logs,
partition,
&before,
)
.await;
return Ok((RepairApplication::AlreadyApplied, 0));
}
if scan.before == before {
let quarantined = scan.quarantined.len();
let effect = apply_repair_scan(
context,
logs,
checkpoints,
count_checkpoints,
mmr_logs,
signer,
partition,
scan,
replacement,
)
.await?;
return Ok((
if effect.rewrote {
RepairApplication::Applied
} else {
RepairApplication::AlreadyApplied
},
quarantined,
));
}
let stage = repair_stage_partition(partition, &before);
let held = {
let log = logs.get_or_open(context, &stage).await?;
log.replay().await?
};
let Some(survivors) = staged_survivors(&held, &before, &after) else {
return Ok((RepairApplication::Changed, 0));
};
let window_open = repair_window_is_open(context, count_checkpoints, partition).await?;
drop(scan);
if !window_open {
return Ok((RepairApplication::Changed, 0));
}
let staged_events = survivors.clone();
let readable = survivors
.into_iter()
.enumerate()
.map(|(position, event)| (position as u64, event))
.collect::<Vec<_>>();
let effect = apply_repair_scan(
context,
logs,
checkpoints,
count_checkpoints,
mmr_logs,
signer,
partition,
RepairScan {
readable,
quarantined: vec![QuarantinedItem {
position: 0,
error: "repair resumed from durable survivor stage".to_owned(),
}],
before,
after,
standing: polyc_eventlog::root_standing_with_trust(&staged_events, trust)?,
},
replacement,
)
.await?;
Ok((
if effect.rewrote {
RepairApplication::Applied
} else {
RepairApplication::AlreadyApplied
},
0,
))
}
#[allow(clippy::too_many_arguments)]
async fn repair_one<E>(
context: &E,
logs: &mut OpenLogs<E>,
checkpoints: &mut HashMap<String, CheckpointFrame>,
count_checkpoints: &mut CountCheckpoints<E>,
mmr_logs: &mut HashMap<String, polyc_mmr::VerifiableLog>,
signer: &JournalAttestationSigner,
trust: &polyc_crypto::signing_role::RoleTrustSet<
polyc_crypto::signing_role::JournalAttestationRole,
>,
partition: &str,
) -> Result<RepairEffect, EventLogError>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
let scan = inspect_repair_one(
context,
logs,
count_checkpoints,
signer,
trust,
partition,
None,
)
.await?;
apply_repair_scan(
context,
logs,
checkpoints,
count_checkpoints,
mmr_logs,
signer,
partition,
scan,
None,
)
.await
}
async fn replay_one<E>(
context: &E,
logs: &mut OpenLogs<E>,
partition: &str,
) -> Result<Vec<Event>, EventLogError>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
#[cfg(test)]
maybe_pause_replay(partition).await;
let log = logs.get_or_open(context, partition).await?;
log.replay().await
}
#[cfg(test)]
async fn maybe_pause_replay(partition: &str) {
if let Some(pause) = tests::take_replay_pause(partition) {
pause.entered.notify_one();
pause.release.notified().await;
}
}
async fn replay_with_positions_one<E>(
context: &E,
logs: &mut OpenLogs<E>,
partition: &str,
) -> Result<Vec<(u64, Event)>, EventLogError>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
let log = logs.get_or_open(context, partition).await?;
log.replay_with_positions().await
}
async fn replay_with_positions_bounded_one<E>(
context: &E,
logs: &mut OpenLogs<E>,
partition: &str,
max_bytes: u64,
) -> Result<BoundedReplay, EventLogError>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
let log = logs.get_or_open(context, partition).await?;
log.replay_with_positions_bounded(max_bytes).await
}
async fn replay_from_with_positions_one<E>(
context: &E,
logs: &mut OpenLogs<E>,
partition: &str,
start: u64,
) -> Result<Vec<(u64, Event)>, EventLogError>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
let log = logs.get_or_open(context, partition).await?;
log.replay_from_with_positions(start).await
}
async fn replay_from_with_positions_bounded_one<E>(
context: &E,
logs: &mut OpenLogs<E>,
partition: &str,
start: u64,
max_bytes: u64,
) -> Result<BoundedReplay, EventLogError>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
let log = logs.get_or_open(context, partition).await?;
log.replay_from_with_positions_bounded(start, max_bytes)
.await
}
async fn replay_range_with_positions_bounded_one<E>(
context: &E,
logs: &mut OpenLogs<E>,
partition: &str,
start: u64,
end: u64,
max_bytes: u64,
) -> Result<BoundedReplay, EventLogError>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
let log = logs.get_or_open(context, partition).await?;
log.replay_range_with_positions_bounded(start, end, max_bytes)
.await
}
async fn events_since_checkpoint_one<E>(
context: &E,
logs: &mut OpenLogs<E>,
checkpoints: &HashMap<String, CheckpointFrame>,
partition: &str,
) -> Result<u64, EventLogError>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
let log = logs.get_or_open(context, partition).await?;
let len = log.len().await;
let since = checkpoints
.get(partition)
.map_or(len, |frame| len.saturating_sub(frame.tail_offset));
Ok(since)
}
async fn partition_event_count_one<E>(
context: &E,
storage_dir: &std::path::Path,
logs: &mut OpenLogs<E>,
partition: &str,
) -> Result<u64, AppendError>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
if !logs.contains(partition) {
if !partition_data_dir_exists(storage_dir, partition)? {
return Ok(0);
}
}
let log = logs
.get_or_open(context, partition)
.await
.map_err(AppendError::Log)?;
Ok(log.len().await)
}
async fn replay_tail_one<E>(
context: &E,
logs: &mut OpenLogs<E>,
checkpoints: &mut HashMap<String, CheckpointFrame>,
partition: &str,
) -> Result<CheckpointReplay, EventLogError>
where
E: commonware_storage::Context + commonware_runtime::BufferPooler + commonware_runtime::Metrics,
{
if let Some(frame) = checkpoints.get(partition) {
let tail_offset = frame.tail_offset;
let payload = frame.payload.clone();
let log = logs.get_or_open(context, partition).await?;
let tail = log.replay_from(tail_offset).await?;
return Ok(CheckpointReplay {
checkpoint: Some(CheckpointFrame {
tail_offset,
payload,
}),
tail,
});
}
let log = logs.get_or_open(context, partition).await?;
let all = log.replay_with_positions().await?;
let latest = all
.iter()
.rev()
.find(|(_, ev)| parse_kind(&ev.kind).0 == BASE_COMPACTION_CHECKPOINT)
.map(|(pos, ev)| (*pos, ev.payload.clone()));
match latest {
Some((pos, payload)) => {
let tail_offset = pos + 1;
checkpoints.insert(
partition.to_owned(),
CheckpointFrame {
tail_offset,
payload: payload.clone(),
},
);
let tail = all
.into_iter()
.filter_map(|(p, ev)| (p >= tail_offset).then_some(ev))
.collect();
Ok(CheckpointReplay {
checkpoint: Some(CheckpointFrame {
tail_offset,
payload,
}),
tail,
})
}
None => Ok(CheckpointReplay {
checkpoint: None,
tail: all.into_iter().map(|(_, ev)| ev).collect(),
}),
}
}
#[cfg(test)]
mod tests {
use super::{
AppendError, CheckpointFrame, EventLogError, EventLogHost, ListPartitionsError,
PartitionMigration, RepairApplication, RewriteDecision, arm_root_append_failure,
};
use polyc_eventlog::Event;
use std::{
collections::{HashMap, HashSet},
sync::{Arc, Mutex, OnceLock},
};
use tokio_util::sync::CancellationToken;
fn shorten_backing_blobs(
dir: &std::path::Path,
partition: &str,
suffix: &str,
raw_len: u64,
) -> usize {
let mut shortened = 0;
let entries = std::fs::read_dir(dir.join(format!("{partition}{suffix}")))
.expect("partition storage directory");
for entry in entries {
let path = entry.expect("directory entry").path();
if !path.is_file() {
continue;
}
let file = std::fs::OpenOptions::new()
.write(true)
.open(&path)
.expect("open the backing blob");
file.set_len(raw_len).expect("shorten below the header");
file.sync_all().expect("make the truncation durable");
shortened += 1;
}
shortened
}
#[tokio::test]
async fn a_data_blob_below_the_runtime_header_reads_as_truncation_not_an_empty_conversation() {
for raw_len in 0..8u64 {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-partial-header-{}-{raw_len}",
std::process::id()
));
let _ = std::fs::remove_dir_all(&dir);
let partition = "conv-partial-header".to_owned();
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(
dir.clone(),
shutdown.clone(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(23),
)
.expect("spawn host");
host.append_batch(
partition.clone(),
vec![
Event::new("user_msg", b"one".to_vec()),
Event::new("output_msg", b"two".to_vec()),
Event::new("tool_call", b"three".to_vec()),
],
)
.await
.expect("seed");
let committed = host
.partition_event_count(partition.clone())
.await
.expect("count");
assert!(committed >= 3, "the batch and its signed root are durable");
drop(host);
shutdown.cancel();
let shortened = shorten_backing_blobs(&dir, &partition, "_data", raw_len);
assert!(shortened > 0, "the case must actually shorten a data blob");
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(
dir.clone(),
shutdown.clone(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(23),
)
.expect("respawn host");
let error = host
.replay_verified(partition.clone())
.await
.expect_err("a shortened data blob must not verify as a short conversation");
match error {
AppendError::Verify(super::VerifyError::TruncatedReplay { expected, actual }) => {
assert_eq!(expected, committed, "the floor still names what committed");
assert_eq!(
actual, 0,
"Commonware reset the blob to empty; the floor is what catches it"
);
}
other => panic!(
"a data blob shortened to {raw_len} byte(s) reported {other} rather than a \
truncated replay"
),
}
assert_eq!(
host.partition_presence(partition.clone())
.await
.expect("presence"),
super::PartitionPresence::Present,
"storage that lost its header is still present; reporting absence would let a \
caller read the conversation as one that never held anything"
);
drop(host);
shutdown.cancel();
let _ = std::fs::remove_dir_all(&dir);
}
}
#[tokio::test]
async fn an_interrupted_reclaim_converges_to_absence_on_the_next_attempt() {
let suffixes = ["_data", "_offsets-blobs", "_offsets-metadata"];
for mask in 0..8u8 {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-partial-reclaim-{}-{mask}",
std::process::id()
));
let _ = std::fs::remove_dir_all(&dir);
let partition = "conv-partial-reclaim".to_owned();
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(
dir.clone(),
shutdown.clone(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(26),
)
.expect("spawn host");
host.append_batch(
partition.clone(),
vec![Event::new("user_msg", b"one".to_vec())],
)
.await
.expect("seed");
drop(host);
shutdown.cancel();
{
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(
dir.clone(),
shutdown.clone(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(26),
)
.expect("spawn host to empty");
host.rewrite_partition_for_test(
partition.clone(),
Box::new(|_: &Event| RewriteDecision::Drop),
)
.await
.expect("empty the partition");
drop(host);
shutdown.cancel();
}
let mut removed = Vec::new();
for (index, suffix) in suffixes.iter().enumerate() {
if mask & (1 << index) == 0 {
continue;
}
let path = dir.join(format!("{partition}{suffix}"));
if path.exists() {
std::fs::remove_dir_all(&path).expect("remove a partition directory");
removed.push(*suffix);
}
}
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(
dir.clone(),
shutdown.clone(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(26),
)
.expect("respawn host");
host.destroy_partition(partition.clone())
.await
.unwrap_or_else(|error| {
panic!(
"a retry after a reclaim that stopped having removed {removed:?} must \
finish rather than strand the erase, got {error}"
)
});
assert_eq!(
host.partition_presence(partition.clone())
.await
.expect("presence"),
super::PartitionPresence::Absent,
"the erase must converge to absence after stopping with {removed:?} removed"
);
drop(host);
shutdown.cancel();
let _ = std::fs::remove_dir_all(&dir);
}
}
#[tokio::test]
async fn one_damaged_event_count_copy_survives_and_two_do_not() {
for (label, damage_all) in [("one copy", false), ("both copies", true)] {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-count-damage-{}-{}",
std::process::id(),
label.replace(' ', "-")
));
let _ = std::fs::remove_dir_all(&dir);
let partition = "conv-count-damage".to_owned();
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(
dir.clone(),
shutdown.clone(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(25),
)
.expect("spawn host");
host.append_batch(
partition.clone(),
vec![
Event::new("user_msg", b"one".to_vec()),
Event::new("output_msg", b"two".to_vec()),
],
)
.await
.expect("seed");
let committed = host
.partition_event_count(partition.clone())
.await
.expect("count");
assert!(committed >= 2, "the batch and its signed root are durable");
drop(host);
shutdown.cancel();
let store = dir.join(format!("{partition}__eventcount_checkpoint"));
let mut copies: Vec<std::path::PathBuf> = std::fs::read_dir(&store)
.expect("checkpoint store directory")
.map(|entry| entry.expect("directory entry").path())
.filter(|path| path.is_file())
.collect();
copies.sort();
assert_eq!(copies.len(), 2, "the metadata store keeps two copies");
let damaged = if damage_all {
&copies[..]
} else {
&copies[..1]
};
for copy in damaged {
let file = std::fs::OpenOptions::new()
.write(true)
.open(copy)
.expect("open a checkpoint copy");
file.set_len(0).expect("shorten below the header");
file.sync_all().expect("make the truncation durable");
}
let shortened = shorten_backing_blobs(&dir, &partition, "_data", 0);
assert!(shortened > 0, "the case must actually shorten a data blob");
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(
dir.clone(),
shutdown.clone(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(25),
)
.expect("respawn host");
let outcome = host.replay_verified(partition.clone()).await;
if damage_all {
let events =
outcome.expect("with no floor left, nothing contradicts the emptied journal");
assert!(
events.is_empty(),
"the records are gone; only the floor could have said so"
);
} else {
match outcome {
Err(AppendError::Verify(super::VerifyError::TruncatedReplay {
expected,
actual,
})) => {
assert_eq!(expected, committed, "the surviving copy holds the floor");
assert_eq!(actual, 0);
}
other => {
panic!("one damaged copy must leave the floor readable, got {other:?}")
}
}
}
drop(host);
shutdown.cancel();
let _ = std::fs::remove_dir_all(&dir);
}
}
#[tokio::test]
async fn an_offsets_blob_below_the_runtime_header_refuses_rather_than_reading_empty() {
for raw_len in 0..8u64 {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-partial-offsets-{}-{raw_len}",
std::process::id()
));
let _ = std::fs::remove_dir_all(&dir);
let partition = "conv-partial-offsets".to_owned();
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(
dir.clone(),
shutdown.clone(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(24),
)
.expect("spawn host");
host.append_batch(
partition.clone(),
vec![Event::new("user_msg", b"one".to_vec())],
)
.await
.expect("seed");
drop(host);
shutdown.cancel();
let shortened = shorten_backing_blobs(&dir, &partition, "_offsets-blobs", raw_len);
assert!(
shortened > 0,
"the case must actually shorten an offsets blob"
);
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(
dir.clone(),
shutdown.clone(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(24),
)
.expect("respawn host");
let outcome = host.replay_verified(partition.clone()).await;
assert!(
outcome.is_err(),
"a shortened offsets blob answered {outcome:?}; an empty history over storage \
Commonware cannot reconcile is the one answer this must never give"
);
drop(host);
shutdown.cancel();
let _ = std::fs::remove_dir_all(&dir);
}
}
#[test]
fn a_failed_checkpoint_sync_evicts_the_handle_and_the_next_attempt_reopens() {
use commonware_runtime::{Runner as _, deterministic, deterministic::FaultConfig};
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let partition = "conv-poisoned-floor";
let mut count_checkpoints = super::CountCheckpoints::new();
count_checkpoints
.record(&context, partition, 5)
.await
.unwrap();
assert!(
count_checkpoints.contains(partition),
"a successful mutation keeps a healthy handle resident"
);
*context.storage_fault_config().write() = FaultConfig::default().sync(1.0);
let error = count_checkpoints
.record(&context, partition, 9)
.await
.expect_err("the durable sync must fail");
assert!(
matches!(error, EventLogError::Checkpoint(_)),
"expected a checkpoint storage error, got {error}"
);
assert!(
!count_checkpoints.contains(partition),
"a failed mutation must leave no handle for a later call to reuse"
);
*context.storage_fault_config().write() = FaultConfig::default();
assert_eq!(
count_checkpoints
.expected_count(&context, partition)
.await
.unwrap(),
5,
"the next attempt reopens from durable storage, so it reads the last count that actually synced rather than the one the failed instance held"
);
});
}
#[test]
fn a_best_effort_floor_failure_still_evicts_the_handle() {
use commonware_runtime::{Runner as _, deterministic, deterministic::FaultConfig};
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let partition = "conv-best-effort-floor";
let mut count_checkpoints = super::CountCheckpoints::new();
count_checkpoints
.record(&context, partition, 3)
.await
.unwrap();
*context.storage_fault_config().write() = FaultConfig::default().sync(1.0);
super::record_event_count_floor(&context, &mut count_checkpoints, partition, 8).await;
assert!(
!count_checkpoints.contains(partition),
"the swallowed failure still evicts the handle"
);
*context.storage_fault_config().write() = FaultConfig::default();
assert_eq!(
count_checkpoints
.expected_count(&context, partition)
.await
.unwrap(),
3,
"the floor stays at its last durable value"
);
});
}
static ARMED_APPEND_FAILURES: OnceLock<Mutex<HashMap<String, usize>>> = OnceLock::new();
static ARMED_COMMIT_FAILURES: OnceLock<Mutex<HashSet<String>>> = OnceLock::new();
static ARMED_MMR_PREPARATION_FAILURES: OnceLock<Mutex<HashSet<String>>> = OnceLock::new();
#[derive(Clone)]
pub(crate) struct ReplayPause {
pub(crate) entered: Arc<tokio::sync::Notify>,
pub(crate) release: Arc<tokio::sync::Notify>,
}
static ARMED_REPLAY_PAUSES: OnceLock<Mutex<HashMap<String, ReplayPause>>> = OnceLock::new();
#[test]
fn exact_rewrite_fingerprints_bind_lineage_payload_and_deterministic_root() {
let signer = polyc_crypto::signing_role::JournalAttestationSigner::from_seed(91);
let events = vec![
Event::new("__state_journal_incarnation__", vec![1; 33]),
Event::new("record", vec![2, 3]),
];
let decision = |event: &Event| {
if event.kind == "__state_journal_incarnation__" {
RewriteDecision::Replace(vec![9; 33])
} else if event.kind == "record" {
RewriteDecision::Replace(vec![4, 5])
} else {
RewriteDecision::Keep
}
};
let (first, restored, _) =
super::inspect_exact_rewrite("conv", "excise", &events, &signer, &decision)
.expect("plan");
let (retry, retry_restored, _) =
super::inspect_exact_rewrite("conv", "excise", &events, &signer, &decision)
.expect("retry plan");
assert_eq!(first, retry, "the same signed result is deterministic");
assert_eq!(restored, retry_restored);
let changed_payload = |event: &Event| {
if event.kind == "__state_journal_incarnation__" {
RewriteDecision::Replace(vec![9; 33])
} else if event.kind == "record" {
RewriteDecision::Replace(vec![4, 6])
} else {
RewriteDecision::Keep
}
};
let (payload_plan, _, _) =
super::inspect_exact_rewrite("conv", "excise", &events, &signer, &changed_payload)
.expect("changed payload plan");
assert_ne!(first.after, payload_plan.after);
assert_ne!(first.stage_binding, payload_plan.stage_binding);
let changed_lineage = |event: &Event| {
if event.kind == "__state_journal_incarnation__" {
RewriteDecision::Replace(vec![8; 33])
} else if event.kind == "record" {
RewriteDecision::Replace(vec![4, 5])
} else {
RewriteDecision::Keep
}
};
let (lineage_plan, _, _) =
super::inspect_exact_rewrite("conv", "excise", &events, &signer, &changed_lineage)
.expect("changed lineage plan");
assert_ne!(first.after, lineage_plan.after);
assert_ne!(first.stage_binding, lineage_plan.stage_binding);
}
#[test]
fn root_append_fault_arms_are_counted_per_partition() {
let partition = "root-fault-count";
arm_root_append_failure(partition);
arm_root_append_failure(partition);
assert!(super::take_armed_root_append_failure(partition).is_some());
assert!(super::take_armed_root_append_failure(partition).is_some());
assert!(super::take_armed_root_append_failure(partition).is_none());
}
fn arm_replay_pause(partition: &str) -> ReplayPause {
let pause = ReplayPause {
entered: Arc::new(tokio::sync::Notify::new()),
release: Arc::new(tokio::sync::Notify::new()),
};
ARMED_REPLAY_PAUSES
.get_or_init(|| Mutex::new(HashMap::new()))
.lock()
.expect("lock poisoned")
.insert(partition.to_owned(), pause.clone());
pause
}
pub(crate) fn take_replay_pause(partition: &str) -> Option<ReplayPause> {
ARMED_REPLAY_PAUSES
.get_or_init(|| Mutex::new(HashMap::new()))
.lock()
.expect("lock poisoned")
.remove(partition)
}
fn injected_error() -> EventLogError {
EventLogError::Journal(commonware_storage::journal::Error::Corruption(
"injected test failure (#1712 fault-injection seam)".to_owned(),
))
}
pub(crate) fn arm_append_failure(partition: &str, event_index: usize) {
ARMED_APPEND_FAILURES
.get_or_init(|| Mutex::new(HashMap::new()))
.lock()
.expect("lock poisoned")
.insert(partition.to_owned(), event_index);
}
pub(crate) fn take_armed_append_failure(
partition: &str,
event_index: usize,
) -> Option<EventLogError> {
let lock = ARMED_APPEND_FAILURES.get()?;
let matched = {
let mut armed = lock.lock().expect("lock poisoned");
if armed.get(partition) == Some(&event_index) {
armed.remove(partition);
true
} else {
false
}
};
matched.then(injected_error)
}
pub(crate) fn arm_commit_failure(partition: &str) {
ARMED_COMMIT_FAILURES
.get_or_init(|| Mutex::new(HashSet::new()))
.lock()
.expect("lock poisoned")
.insert(partition.to_owned());
}
pub(crate) fn take_armed_commit_failure(partition: &str) -> Option<EventLogError> {
let lock = ARMED_COMMIT_FAILURES.get()?;
let matched = lock.lock().expect("lock poisoned").remove(partition);
matched.then(injected_error)
}
static ARMED_COUNT_RECORD_FAILURES: OnceLock<Mutex<HashSet<String>>> = OnceLock::new();
pub(crate) fn arm_count_record_failure(partition: &str) {
ARMED_COUNT_RECORD_FAILURES
.get_or_init(|| Mutex::new(HashSet::new()))
.lock()
.expect("lock poisoned")
.insert(partition.to_owned());
}
pub(crate) fn take_armed_count_record_failure(partition: &str) -> Option<EventLogError> {
let lock = ARMED_COUNT_RECORD_FAILURES.get()?;
let matched = lock.lock().expect("lock poisoned").remove(partition);
matched.then(injected_error)
}
fn arm_mmr_preparation_failure(partition: &str) {
ARMED_MMR_PREPARATION_FAILURES
.get_or_init(|| Mutex::new(HashSet::new()))
.lock()
.expect("lock poisoned")
.insert(partition.to_owned());
}
pub(crate) fn take_armed_mmr_preparation_failure(partition: &str) -> Option<EventLogError> {
let lock = ARMED_MMR_PREPARATION_FAILURES.get()?;
let matched = lock.lock().expect("lock poisoned").remove(partition);
matched.then(injected_error)
}
static ARMED_DESTROY_FAILURES: OnceLock<Mutex<HashSet<String>>> = OnceLock::new();
pub(crate) fn arm_destroy_failure(partition: &str) {
ARMED_DESTROY_FAILURES
.get_or_init(|| Mutex::new(HashSet::new()))
.lock()
.expect("lock poisoned")
.insert(partition.to_owned());
}
pub(crate) fn take_armed_destroy_failure(partition: &str) -> Option<EventLogError> {
let lock = ARMED_DESTROY_FAILURES.get()?;
let matched = lock.lock().expect("lock poisoned").remove(partition);
matched.then(injected_error)
}
#[tokio::test]
async fn replay_with_positions_bounded_stops_early_through_the_host() {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-bounded-host-{}",
std::process::id()
));
let _ = std::fs::remove_dir_all(&dir);
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(
dir.clone(),
shutdown.clone(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
)
.expect("spawn host");
let events: Vec<Event> = (0..10)
.map(|i| Event::new(format!("k{i}"), vec![0u8; 1_000]))
.collect();
host.append_batch("conv-bounded-host".to_owned(), events)
.await
.expect("seed");
let bounded = host
.replay_with_positions_bounded("conv-bounded-host".to_owned(), 3_500)
.await
.expect("bounded replay");
assert!(
bounded.budget_exceeded,
"a 3,500 byte budget over 10,000 bytes of payload must trip"
);
assert!(
bounded.events.len() < 10,
"the host must not hand back every event once the budget trips: got {}",
bounded.events.len()
);
assert_eq!(
bounded.events.len(),
4,
"stops right after the 4th 1,000-byte event"
);
let unbounded = host
.replay_with_positions("conv-bounded-host".to_owned())
.await
.expect("unbounded replay");
assert!(
bounded.events.len() < unbounded.len(),
"the bounded replay must return strictly fewer events than the unbounded one: \
bounded={} unbounded={}",
bounded.events.len(),
unbounded.len()
);
drop(host);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn replay_range_with_positions_bounded_through_the_host() {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-range-bounded-host-{}",
std::process::id()
));
let _ = std::fs::remove_dir_all(&dir);
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(
dir.clone(),
shutdown.clone(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
)
.expect("spawn host");
let events: Vec<Event> = (0..10)
.map(|i| Event::new(format!("k{i}"), vec![0u8; 1_000]))
.collect();
host.append_batch("conv-range-bounded-host".to_owned(), events)
.await
.expect("seed");
let ranged = host
.replay_range_with_positions_bounded(
"conv-range-bounded-host".to_owned(),
3,
6,
u64::MAX,
)
.await
.expect("range replay");
let positions: Vec<u64> = ranged.events.iter().map(|(p, _)| *p).collect();
assert_eq!(
positions,
vec![3, 4, 5],
"[3, 6) is 3, 4, 5 — position 6 is outside the range"
);
assert!(
!ranged.budget_exceeded,
"a complete range must not report itself truncated"
);
assert_eq!(ranged.bytes_read, 3_000);
let capped = host
.replay_range_with_positions_bounded("conv-range-bounded-host".to_owned(), 2, 9, 2_500)
.await
.expect("range replay under a budget");
assert!(capped.budget_exceeded);
let capped_positions: Vec<u64> = capped.events.iter().map(|(p, _)| *p).collect();
assert_eq!(capped_positions, vec![2, 3, 4]);
let clamped = host
.replay_range_with_positions_bounded(
"conv-range-bounded-host".to_owned(),
8,
9_999,
u64::MAX,
)
.await
.expect("range replay past the tail");
let clamped_positions: Vec<u64> = clamped.events.iter().map(|(p, _)| *p).collect();
assert_eq!(
clamped_positions,
vec![8, 9, 10],
"clamped to the REAL tail, which is 11 events rather than the 10 seeded: \
`append_batch` also writes the host's own per-commit `__mmr_signed_root__` \
tamper-evidence marker at position 10. A range clamp has to land on what the \
journal actually holds, not on what the caller appended"
);
let empty = host
.replay_range_with_positions_bounded("conv-range-bounded-host".to_owned(), 5, 5, 1)
.await
.expect("empty range");
assert!(empty.events.is_empty());
assert!(!empty.budget_exceeded);
drop(host);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn replay_from_with_positions_bounded_honors_both_through_the_host() {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-tail-bounded-host-{}",
std::process::id()
));
let _ = std::fs::remove_dir_all(&dir);
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(
dir.clone(),
shutdown.clone(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
)
.expect("spawn host");
let events: Vec<Event> = (0..10)
.map(|i| Event::new(format!("k{i}"), vec![0u8; 1_000]))
.collect();
host.append_batch("conv-tail-bounded-host".to_owned(), events)
.await
.expect("seed");
let bounded = host
.replay_from_with_positions_bounded("conv-tail-bounded-host".to_owned(), 4, 2_500)
.await
.expect("bounded tail replay");
assert!(
bounded.budget_exceeded,
"a 2,500 byte budget over the 6,000-byte tail (positions 4..10) must trip"
);
let positions: Vec<u64> = bounded.events.iter().map(|(p, _)| *p).collect();
assert_eq!(
positions,
vec![4, 5, 6],
"must resume at position 4 (never re-reading 0..4) and stop the instant \
cumulative bytes (3,000 after the 3rd tail event) cross the 2,500 budget"
);
assert_eq!(bounded.bytes_read, 3_000);
let past_end = host
.replay_from_with_positions_bounded("conv-tail-bounded-host".to_owned(), 999, 1)
.await
.expect("bounded tail replay past the end");
assert!(past_end.events.is_empty());
assert!(!past_end.budget_exceeded);
drop(host);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn list_partitions_returns_logical_names_and_tracks_destroy() {
let dir =
std::env::temp_dir().join(format!("polychrome-eventlog-list-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(
dir.clone(),
shutdown.clone(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
)
.expect("spawn host");
assert!(
host.list_partitions().await.expect("list empty").is_empty(),
"a fresh host has no partitions"
);
for partition in ["conv-c1", "persona-p1-mem", "erasure-audit"] {
host.append_batch(
partition.to_owned(),
vec![Event::new("output_msg", b"x".to_vec())],
)
.await
.expect("seed");
}
assert_eq!(
host.list_partitions().await.expect("list"),
vec!["conv-c1", "erasure-audit", "persona-p1-mem"],
"logical names, sorted — no storage-layout suffixes"
);
assert!(matches!(
host.list_partitions_bounded(2, usize::MAX).await,
Err(AppendError::Listing(ListPartitionsError::EntriesExceeded {
limit: 2,
requested: 3,
}))
));
assert!(matches!(
host.list_partitions_bounded(3, "conv-c1".len()).await,
Err(AppendError::Listing(
ListPartitionsError::NameBytesExceeded { .. }
))
));
host.destroy_partition("persona-p1-mem".to_owned())
.await
.expect("destroy");
assert_eq!(
host.list_partitions().await.expect("list after destroy"),
vec!["conv-c1", "erasure-audit"],
"a destroyed partition no longer lists"
);
drop(host);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn a_corrupt_partition_shaped_entry_refuses_the_whole_listing() {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-corrupt-list-{}",
std::process::id()
));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(dir.join("good_data")).expect("valid partition directory");
std::fs::create_dir_all(dir.join("_61_data")).expect("corrupt partition directory");
assert!(matches!(
crate::list_partitions_bounded_in(&dir, usize::MAX, usize::MAX, None),
Err(ListPartitionsError::Corrupt { entry }) if entry == "_61_data"
));
}
#[test]
fn a_retry_after_a_failed_marker_append_still_re_roots() {
use commonware_runtime::{Runner as _, tokio as cw_tokio};
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-marker-retry-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let _ = std::fs::remove_dir_all(&dir);
let to = "conv-marker-retry-to".to_owned();
let from = "conv-marker-retry-from".to_owned();
let to_ref = super::PartitionRef::new(&to).expect("the name encodes");
let from_ref = super::PartitionRef::new(&from).expect("the name encodes");
let runner =
cw_tokio::Runner::new(cw_tokio::Config::default().with_storage_directory(dir.clone()));
runner.start(move |context| async move {
let mut logs = super::OpenLogs::new();
let mut count_checkpoints = super::CountCheckpoints::new();
let mut mmr_logs = HashMap::new();
let write_observers = super::observer::WriteObservers::new();
let signer = polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1);
let candidates = [
Event::new("memory_added", b"one".to_vec()),
Event::new("memory_added", b"two".to_vec()),
];
arm_append_failure(&to, 2);
super::migrate_append_into(
&context,
&mut logs,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&to_ref,
&candidates,
&write_observers,
&from_ref,
)
.await
.expect_err("the armed marker append fails the first attempt");
let appended = super::migrate_append_into(
&context,
&mut logs,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&to_ref,
&candidates,
&write_observers,
&from_ref,
)
.await
.expect("the retry succeeds");
assert_eq!(appended, 0, "nothing new moved on the retry");
let log = logs.get_or_open(&context, &to).await.expect("reopen");
let replayed = log.replay().await.expect("replay");
assert!(
replayed
.iter()
.any(|ev| ev.kind == polyc_eventlog::MMR_SIGNED_ROOT_KIND),
"the retry must leave a signed root over the migrated events"
);
assert_eq!(
replayed.last().map(|ev| ev.kind.as_str()),
Some(polyc_eventlog::MMR_SIGNED_ROOT_KIND),
"and it must cover everything, so it is last"
);
});
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn a_retry_re_roots_a_destination_that_already_holds_a_root() {
use commonware_runtime::{Runner as _, tokio as cw_tokio};
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-marker-retry-nonempty-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let _ = std::fs::remove_dir_all(&dir);
let to = "conv-marker-retry-nonempty-to".to_owned();
let from = "conv-marker-retry-nonempty-from".to_owned();
let to_ref = super::PartitionRef::new(&to).expect("the name encodes");
let from_ref = super::PartitionRef::new(&from).expect("the name encodes");
let runner =
cw_tokio::Runner::new(cw_tokio::Config::default().with_storage_directory(dir.clone()));
runner.start(move |context| async move {
let mut logs = super::OpenLogs::new();
let mut count_checkpoints = super::CountCheckpoints::new();
let mut mmr_logs = HashMap::new();
let write_observers = super::observer::WriteObservers::new();
let signer = polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1);
let seed = [Event::new("memory_added", b"seed".to_vec())];
let candidates = [
Event::new("memory_added", b"one".to_vec()),
Event::new("memory_added", b"two".to_vec()),
];
super::migrate_append_into(
&context,
&mut logs,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&to_ref,
&seed,
&write_observers,
&from_ref,
)
.await
.expect("the seeding migration succeeds");
arm_append_failure(&to, 2);
super::migrate_append_into(
&context,
&mut logs,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&to_ref,
&candidates,
&write_observers,
&from_ref,
)
.await
.expect_err("the armed marker append fails the first attempt");
let log = logs.get_or_open(&context, &to).await.expect("reopen");
let mid = log.replay().await.expect("replay");
assert_eq!(
mid.last().map(|ev| ev.kind.as_str()),
Some("memory_added"),
"the failed attempt must leave the tail unattested, or this \
test proves nothing about the retry"
);
let appended = super::migrate_append_into(
&context,
&mut logs,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&to_ref,
&candidates,
&write_observers,
&from_ref,
)
.await
.expect("the retry succeeds");
assert_eq!(appended, 0, "nothing new moved on the retry");
let log = logs.get_or_open(&context, &to).await.expect("reopen");
let replayed = log.replay().await.expect("replay");
assert_eq!(
replayed.last().map(|ev| ev.kind.as_str()),
Some(polyc_eventlog::MMR_SIGNED_ROOT_KIND),
"the retry must re-root, and the root must be last"
);
assert_eq!(
replayed
.iter()
.filter(|ev| ev.kind == polyc_eventlog::MMR_SIGNED_ROOT_KIND)
.count(),
2,
"the destination's own root survives the re-root"
);
let trust = polyc_crypto::signing_role::RoleTrustSet::current(&signer);
polyc_eventlog::verify_replay_with_trust(&replayed, &trust)
.expect("both roots verify: R0 over one leaf, R1 over three");
});
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn a_re_root_cannot_repair_a_root_that_covers_the_wrong_leaves() {
use commonware_runtime::{Runner as _, tokio as cw_tokio};
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-foreign-root-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let _ = std::fs::remove_dir_all(&dir);
let to = "conv-foreign-root-to".to_owned();
let from = "conv-foreign-root-from".to_owned();
let to_ref = super::PartitionRef::new(&to).expect("the name encodes");
let from_ref = super::PartitionRef::new(&from).expect("the name encodes");
let runner =
cw_tokio::Runner::new(cw_tokio::Config::default().with_storage_directory(dir.clone()));
runner.start(move |context| async move {
let mut logs = super::OpenLogs::new();
let mut count_checkpoints = super::CountCheckpoints::new();
let mut mmr_logs = HashMap::new();
let write_observers = super::observer::WriteObservers::new();
let signer = polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1);
let carried = [
Event::new("memory_added", b"one".to_vec()),
Event::new("memory_added", b"two".to_vec()),
Event::new("memory_added", b"three".to_vec()),
];
{
let log = logs.get_or_open(&context, &to).await.expect("open");
for event in &carried {
log.append(event).await.expect("append");
}
let short = polyc_eventlog::rebuild_from_events(&carried[..1]).expect("rebuild");
let foreign = polyc_eventlog::extend_and_sign(&short, &[], &signer).expect("sign");
log.append(&foreign).await.expect("append the foreign root");
log.commit().await.expect("commit");
}
let trust = polyc_crypto::signing_role::RoleTrustSet::current(&signer);
{
let log = logs.get_or_open(&context, &to).await.expect("reopen");
let poisoned = log.replay().await.expect("replay");
assert_eq!(
poisoned.last().map(|ev| ev.kind.as_str()),
Some(polyc_eventlog::MMR_SIGNED_ROOT_KIND),
"the setup must end in a marker, or it does not model the case"
);
assert!(
polyc_eventlog::verify_replay_with_trust(&poisoned, &trust).is_err(),
"the setup must actually be poisoned, or this test proves nothing"
);
}
let appended = super::migrate_append_into(
&context,
&mut logs,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&to_ref,
&carried,
&write_observers,
&from_ref,
)
.await
.expect("the migration succeeds");
assert_eq!(appended, 0, "nothing new moved");
let log = logs.get_or_open(&context, &to).await.expect("reopen");
let replayed = log.replay().await.expect("replay");
let still_broken = polyc_eventlog::verify_replay_with_trust(&replayed, &trust)
.expect_err("a poisoned partition stays poisoned; this fix is forward-only");
assert!(
matches!(
still_broken,
polyc_eventlog::IntegrityError::RootMismatch { leaf_count: 3, .. }
),
"it must fail on the bad root it already holds, over the three \
leaves present, not on something new: {still_broken:?}"
);
});
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn a_concurrent_append_cannot_poison_a_migration() {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-migrate-race-{}",
std::process::id()
));
let _ = std::fs::remove_dir_all(&dir);
let shutdown = CancellationToken::new();
let host = std::sync::Arc::new(
EventLogHost::spawn(
dir.clone(),
shutdown.clone(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
)
.expect("spawn host"),
);
host.append_batch(
"conv-race-dest".to_owned(),
vec![Event::new("output_msg", b"dest-one".to_vec())],
)
.await
.expect("seed destination");
let bulk: Vec<Event> = (0..2000)
.map(|i| Event::new("output_msg", format!("src-{i}").into_bytes()))
.collect();
host.append_batch("conv-race-source".to_owned(), bulk)
.await
.expect("seed source");
let racer = {
let host = std::sync::Arc::clone(&host);
tokio::spawn(async move {
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
host.append_batch(
"conv-race-dest".to_owned(),
vec![Event::new("output_msg", b"concurrent".to_vec())],
)
.await
})
};
host.migrate_partition("conv-race-source".to_owned(), "conv-race-dest".to_owned())
.await
.expect("migrate");
racer
.await
.expect("the racing task ran")
.expect("racing append");
host.replay_verified("conv-race-dest".to_owned())
.await
.expect("a concurrent append must not poison the destination");
shutdown.cancel();
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn a_migrated_partition_verifies_at_its_destination() {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-migrate-reroot-{}",
std::process::id()
));
let _ = std::fs::remove_dir_all(&dir);
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(
dir.clone(),
shutdown.clone(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
)
.expect("spawn host");
host.append_batch(
"conv-migrate-source".to_owned(),
vec![Event::new("output_msg", b"source-one".to_vec())],
)
.await
.expect("seed source");
host.append_batch(
"conv-migrate-dest".to_owned(),
vec![Event::new("output_msg", b"dest-one".to_vec())],
)
.await
.expect("seed destination");
host.replay_verified("conv-migrate-dest".to_owned())
.await
.expect("the destination verifies before the migration");
host.migrate_partition(
"conv-migrate-source".to_owned(),
"conv-migrate-dest".to_owned(),
)
.await
.expect("migrate");
let verified = host
.replay_verified("conv-migrate-dest".to_owned())
.await
.expect("the destination must still verify after the migration");
let payloads: Vec<&[u8]> = verified.iter().map(|ev| ev.payload.as_slice()).collect();
assert!(
payloads.contains(&b"source-one".as_slice()),
"the migrated event must arrive"
);
assert!(
payloads.contains(&b"dest-one".as_slice()),
"the destination's own event must survive"
);
let last = verified.last().expect("the destination is not empty");
assert_eq!(
last.kind,
polyc_eventlog::MMR_SIGNED_ROOT_KIND,
"a migration must re-root its destination over what it now holds"
);
shutdown.cancel();
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
#[allow(clippy::too_many_lines)]
async fn migrate_partition_moves_dedups_and_destroys() {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-migrate-{}",
std::process::id()
));
let _ = std::fs::remove_dir_all(&dir);
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(
dir.clone(),
shutdown.clone(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
)
.expect("spawn host");
let a = Event::new("memory_added", b"fact-a".to_vec());
let b = Event::new("memory_added", b"fact-b".to_vec());
let s = Event::new("memory_added", b"fact-s".to_vec());
host.append_batch(
"persona-absorbed-mem".to_owned(),
vec![a.clone(), b.clone()],
)
.await
.expect("seed source");
host.append_batch("persona-survivor-mem".to_owned(), vec![s])
.await
.expect("seed destination");
let outcome = host
.migrate_partition(
"persona-absorbed-mem".to_owned(),
"persona-survivor-mem".to_owned(),
)
.await
.expect("migrate");
assert_eq!(outcome, PartitionMigration::Migrated(2));
assert_eq!(
host.replay("persona-survivor-mem".to_owned())
.await
.expect("replay survivor")
.len(),
5,
"survivor holds its own event + marker, the two moved events, and \
one fresh root over the result"
);
assert!(
host.replay("persona-absorbed-mem".to_owned())
.await
.expect("replay source")
.is_empty(),
"the source is destroyed"
);
host.append_batch("persona-absorbed-mem".to_owned(), vec![a, b])
.await
.expect("re-seed source");
let again = host
.migrate_partition(
"persona-absorbed-mem".to_owned(),
"persona-survivor-mem".to_owned(),
)
.await
.expect("re-migrate");
assert_eq!(
again,
PartitionMigration::Migrated(0),
"dedup: no re-append"
);
assert_eq!(
host.replay("persona-survivor-mem".to_owned())
.await
.expect("replay survivor again")
.len(),
5,
"the survivor journal did not grow"
);
host.append_batch(
"persona-empty-mem".to_owned(),
vec![Event::new("memory_added", b"x".to_vec())],
)
.await
.expect("seed empty candidate");
host.rewrite_partition_for_test(
"persona-empty-mem".to_owned(),
Box::new(|_| RewriteDecision::Drop),
)
.await
.expect("empty it out");
assert_eq!(
host.migrate_partition(
"persona-empty-mem".to_owned(),
"persona-survivor-mem".to_owned()
)
.await
.expect("migrate empty"),
PartitionMigration::EmptyDestroyed
);
assert_eq!(
host.migrate_partition(
"persona-survivor-mem".to_owned(),
"persona-survivor-mem".to_owned()
)
.await
.expect("self-migrate"),
PartitionMigration::Migrated(0)
);
assert_eq!(
host.replay("persona-survivor-mem".to_owned())
.await
.expect("replay after self-migrate")
.len(),
5,
"a self-move destroys nothing"
);
drop(host);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn rewrite_partition_replaces_a_payload_in_place() {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-rewrite-replace-{}",
std::process::id()
));
let _ = std::fs::remove_dir_all(&dir);
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(
dir.clone(),
shutdown.clone(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(3),
)
.expect("spawn host");
host.append_batch(
"persona-rw-mem".to_owned(),
vec![
Event::new("memory_added", b"keep".to_vec()),
Event::new("memory_corroborated", b"has conv-2 and conv-3".to_vec()),
Event::new("memory_added", b"drop-me".to_vec()),
],
)
.await
.expect("seed");
let dropped = host
.rewrite_partition_for_test(
"persona-rw-mem".to_owned(),
Box::new(|e: &Event| {
if e.payload == b"drop-me" {
RewriteDecision::Drop
} else if e.payload.starts_with(b"has ") {
RewriteDecision::Replace(b"has conv-3".to_vec())
} else {
RewriteDecision::Keep
}
}),
)
.await
.expect("rewrite");
assert_eq!(dropped, 1, "a replaced event is not counted as dropped");
let content = |raw: &[Event]| -> Vec<Vec<u8>> {
raw.iter()
.filter(|e| e.kind != "__mmr_signed_root__")
.map(|e| e.payload.clone())
.collect()
};
let raw = host
.replay("persona-rw-mem".to_owned())
.await
.expect("replay");
let payloads = content(&raw);
assert_eq!(
payloads.len(),
2,
"one dropped, one replaced-in-place, one kept: {payloads:?}"
);
assert!(payloads.iter().any(|p| p == b"keep"));
assert!(
payloads.iter().any(|p| p == b"has conv-3"),
"the replaced payload is the new bytes, with conv-2 scrubbed"
);
assert!(
!payloads.iter().any(|p| p == b"has conv-2 and conv-3"),
"the old payload is physically gone"
);
let dropped = host
.rewrite_partition_for_test(
"persona-rw-mem".to_owned(),
Box::new(|e: &Event| {
if e.payload == b"has conv-3" {
RewriteDecision::Replace(b"has nothing".to_vec())
} else {
RewriteDecision::Keep
}
}),
)
.await
.expect("replace-only rewrite");
assert_eq!(dropped, 0, "replace-only rewrite drops nothing");
let raw = host
.replay("persona-rw-mem".to_owned())
.await
.expect("replay again");
assert!(
content(&raw).iter().any(|p| p == b"has nothing"),
"the replace-only rewrite took effect"
);
drop(host);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn rewrite_partition_preserves_each_event_trust_tag() {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-rewrite-trust-{}",
std::process::id()
));
let _ = std::fs::remove_dir_all(&dir);
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(
dir.clone(),
shutdown.clone(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(11),
)
.expect("spawn host");
let partition = "persona-trust-mem".to_owned();
host.append_batch(
partition.clone(),
vec![
Event::trusted("user_msg", b"a".to_vec()),
Event::quarantined("tool_result", b"has conv-2 and conv-3".to_vec()),
],
)
.await
.expect("seed");
let dropped = host
.rewrite_partition_for_test(
partition.clone(),
Box::new(|e: &Event| {
if e.payload.starts_with(b"has ") {
RewriteDecision::Replace(b"has conv-3".to_vec())
} else {
RewriteDecision::Keep
}
}),
)
.await
.expect("scrub");
assert_eq!(dropped, 0);
let survivors: Vec<(String, polyc_eventlog::TrustTag, Vec<u8>)> = host
.replay(partition.clone())
.await
.expect("replay")
.into_iter()
.filter(|e| e.kind != polyc_eventlog::MMR_SIGNED_ROOT_KIND)
.map(|e| (e.kind, e.trust, e.payload))
.collect();
assert_eq!(
survivors,
vec![
(
"user_msg".to_owned(),
polyc_eventlog::TrustTag::TrustedUser,
b"a".to_vec()
),
(
"tool_result".to_owned(),
polyc_eventlog::TrustTag::QuarantinedContent,
b"has conv-3".to_vec()
),
],
"the scrubbed event keeps its kind and its quarantine, and only its \
payload changed"
);
drop(host);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn an_identity_rewrite_over_tagged_events_is_a_no_op() {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-rewrite-identity-{}",
std::process::id()
));
let _ = std::fs::remove_dir_all(&dir);
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(
dir.clone(),
shutdown.clone(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(12),
)
.expect("spawn host");
let partition = "persona-identity-mem".to_owned();
host.append_batch(
partition.clone(),
vec![
Event::trusted("user_msg", b"a".to_vec()),
Event::quarantined("tool_result", b"b".to_vec()),
],
)
.await
.expect("seed");
let roots = |raw: &[Event]| -> Vec<Vec<u8>> {
raw.iter()
.filter(|e| e.kind == polyc_eventlog::MMR_SIGNED_ROOT_KIND)
.map(|e| e.payload.clone())
.collect()
};
let before = host.replay(partition.clone()).await.expect("replay before");
assert_eq!(roots(&before).len(), 1, "the seed committed one root");
let dropped = host
.rewrite_partition_for_test(
partition.clone(),
Box::new(|e: &Event| RewriteDecision::Replace(e.payload.clone())),
)
.await
.expect("identity rewrite");
assert_eq!(dropped, 0);
let after = host.replay(partition.clone()).await.expect("replay after");
assert_eq!(
after, before,
"an identity rewrite leaves the partition byte-for-byte as it was — \
same events, same tags, same single root at the same position"
);
drop(host);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn destroy_partition_evicts_the_handle_and_erases() {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-destroy-{}",
std::process::id()
));
let _ = std::fs::remove_dir_all(&dir);
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(
dir.clone(),
shutdown.clone(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
)
.expect("spawn host");
let partition = "persona-p1-mem".to_owned();
host.append_batch(
partition.clone(),
vec![
Event::new("memory_added", b"secret fact".to_vec()),
Event::new("memory_added", b"another".to_vec()),
],
)
.await
.expect("seed");
assert_eq!(
host.replay(partition.clone()).await.expect("replay").len(),
3
);
host.destroy_partition(partition.clone())
.await
.expect("destroy");
assert!(
host.replay(partition.clone())
.await
.expect("replay after destroy")
.is_empty(),
"a destroyed partition replays empty"
);
let positions = host
.append_batch(
partition.clone(),
vec![Event::new("memory_added", b"fresh".to_vec())],
)
.await
.expect("append after destroy");
assert_eq!(positions, vec![0], "the fresh partition starts at zero");
let events = host.replay(partition).await.expect("replay fresh");
assert_eq!(
events.len(),
2,
"the content event plus its signed-root marker"
);
assert_eq!(events[0].payload, b"fresh", "no old event survives");
drop(host);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn append_batch_and_replay_over_the_bridge() {
let dir =
std::env::temp_dir().join(format!("polychrome-eventlog-test-{}", std::process::id()));
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(
dir.clone(),
shutdown.clone(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
)
.expect("spawn host");
let partition = "conv-test-e1".to_owned();
let positions = host
.append_batch(
partition.clone(),
vec![
Event::new("output_msg", b"one".to_vec()),
Event::new("output_msg", b"two".to_vec()),
],
)
.await
.expect("batch append");
assert_eq!(positions, vec![0, 1], "positions in append order");
let events = host.replay(partition).await.expect("replay");
assert_eq!(
events.len(),
3,
"both content events plus the batch's signed-root marker (#799)"
);
assert_eq!(events[0].payload, b"one");
assert_eq!(events[1].payload, b"two");
let other = host
.replay("conv-other".to_owned())
.await
.expect("replay other");
assert!(other.is_empty(), "untouched partition replays empty");
drop(host);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn two_host_shared_dir_cold_reader_never_sees_another_writers_events() {
fn content(events: &[Event]) -> usize {
events.iter().filter(|e| e.kind == "output_msg").count()
}
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-818-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let _ = std::fs::remove_dir_all(&dir);
let signer = polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1);
let host_a = EventLogHost::spawn(dir.clone(), CancellationToken::new(), signer.clone())
.expect("spawn A");
let host_b = EventLogHost::spawn(dir.clone(), CancellationToken::new(), signer.clone())
.expect("spawn B");
let partition = "conv-818".to_owned();
assert_eq!(
content(
&host_b
.replay(partition.clone())
.await
.expect("B cold replay")
),
0,
"the conversation has no events yet"
);
host_a
.append_batch(
partition.clone(),
vec![
Event::new("output_msg", b"one".to_vec()),
Event::new("output_msg", b"two".to_vec()),
],
)
.await
.expect("A append");
assert_eq!(
content(&host_b.replay(partition.clone()).await.expect("B replay")),
0,
"#818: B's frozen handle never observes A's writes — split-brain \
`raw_events: 0` on a non-writer replica"
);
let host_c =
EventLogHost::spawn(dir.clone(), CancellationToken::new(), signer).expect("spawn C");
assert_eq!(
content(&host_c.replay(partition.clone()).await.expect("C replay")),
2,
"a fresh open reads the full on-disk history — every write is durable"
);
drop(host_a);
drop(host_b);
drop(host_c);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn namespaced_partition_round_trips() {
let dir =
std::env::temp_dir().join(format!("polychrome-eventlog-ns-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(
dir.clone(),
shutdown.clone(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
)
.expect("spawn host");
let partition = "conv-web:019f0a8a-0aca-7410-b1bf-9f644ea1bfe8".to_owned();
host.append_batch(
partition.clone(),
vec![Event::new("output_msg", b"hi".to_vec())],
)
.await
.expect("append to a namespaced partition must succeed");
let events = host
.replay(partition)
.await
.expect("replay of a namespaced partition must succeed");
assert_eq!(
events.len(),
2,
"the namespaced conversation persisted, plus its signed-root marker (#799)"
);
assert_eq!(events[0].payload, b"hi");
drop(host);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn bounded_replay_of_fresh_namespaced_partition_is_empty_not_an_error() {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-ns-tail-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let _ = std::fs::remove_dir_all(&dir);
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(
dir.clone(),
shutdown,
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
)
.expect("spawn host");
let replay = host
.replay_tail_from_checkpoint("conv-web:019f0a8a-0aca-7410-b1bf-9f644ea1bfe8".to_owned())
.await
.expect("bounded replay of a fresh namespaced partition must succeed");
assert!(replay.checkpoint.is_none(), "no checkpoint on a fresh log");
assert!(replay.tail.is_empty(), "a never-written partition is empty");
drop(host);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn append_batch_surfaces_error_when_storage_dir_is_unwritable() {
const PARTITION: &str = "conv-test-bad-dir";
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-blocked-{}",
std::process::id()
));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).expect("create root");
std::fs::write(
dir.join(format!("{PARTITION}_data")),
b"blocks the partition dir",
)
.expect("write blocking file");
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(
dir.clone(),
shutdown,
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
)
.expect("spawn ok (open is lazy)");
let err = host
.append_batch(
PARTITION.to_owned(),
vec![Event::new("output_msg", b"never persisted".to_vec())],
)
.await
.expect_err("append must fail when storage is unusable");
match err {
AppendError::Log(_) => {}
AppendError::Closed
| AppendError::AppendOutcomeUnknown(_)
| AppendError::AttestationOutcomeUnknown(_)
| AppendError::Listing(_)
| AppendError::PartitionName(_)
| AppendError::PayloadTooLarge { .. }
| AppendError::Verify(_) => {
panic!("expected a Log error, got: {err}")
}
}
drop(host);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn ensuring_an_absent_partition_creates_nothing() {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-ensure-absent-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let host = EventLogHost::spawn(
dir.clone(),
CancellationToken::new(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
)
.expect("spawn host");
assert!(
host.ensure_partition_attested("conv-absent".to_owned())
.await
.expect("absence is an answered state")
.is_none()
);
assert!(
!host
.list_partitions()
.await
.expect("list after absence probe")
.contains(&"conv-absent".to_owned())
);
drop(host);
let _ = std::fs::remove_dir_all(dir);
}
#[tokio::test]
async fn mmr_preparation_failure_is_definite_and_appends_no_content() {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-mmr-prepare-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let partition = "conv-mmr-prepare";
let host = EventLogHost::spawn(
dir.clone(),
CancellationToken::new(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
)
.expect("spawn host");
arm_mmr_preparation_failure(partition);
let error = host
.append_batch(
partition.to_owned(),
vec![Event::new("output_msg", b"must-not-land".to_vec())],
)
.await
.expect_err("preparation failure is reported before append");
assert!(
matches!(error, AppendError::Log(_)),
"pre-storage preparation is a definite failure: {error}"
);
assert!(
host.replay(partition.to_owned())
.await
.expect("empty replay")
.is_empty(),
"neither content nor a root marker may land"
);
host.append_batch(
partition.to_owned(),
vec![Event::new("output_msg", b"retry".to_vec())],
)
.await
.expect("a clean retry succeeds");
host.replay_verified(partition.to_owned())
.await
.expect("retry evidence verifies");
drop(host);
let _ = std::fs::remove_dir_all(dir);
}
#[tokio::test]
async fn root_append_ambiguity_retries_to_exactly_one_attestation() {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-root-recovery-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let partition = "conv-root-recovery";
let host = EventLogHost::spawn(
dir.clone(),
CancellationToken::new(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
)
.expect("spawn host");
arm_root_append_failure(partition);
let error = host
.append_batch(
partition.to_owned(),
vec![Event::new("output_msg", b"durable content".to_vec())],
)
.await
.expect_err("root failure leaves the command outcome ambiguous");
assert!(matches!(error, AppendError::AttestationOutcomeUnknown(_)));
arm_root_append_failure(partition);
let ensure_error = host
.ensure_partition_attested(partition.to_owned())
.await
.expect_err("an uncertain ensure append remains ambiguous");
assert!(matches!(
ensure_error,
AppendError::AttestationOutcomeUnknown(_)
));
let recovered = host
.ensure_partition_attested(partition.to_owned())
.await
.expect("reopen and recovery succeed")
.expect("the non-empty prefix receives evidence");
let replayed = host.replay(partition.to_owned()).await.expect("replay");
assert!(
replayed
.iter()
.any(|event| event.payload == b"durable content")
);
assert_eq!(
replayed
.iter()
.filter(|event| event.kind == polyc_eventlog::MMR_SIGNED_ROOT_KIND)
.count(),
1
);
let repeated = host
.ensure_partition_attested(partition.to_owned())
.await
.expect("repeat recovery succeeds")
.expect("existing evidence is returned");
assert_eq!(repeated.kind, recovered.kind);
assert_eq!(repeated.payload, recovered.payload);
host.replay_verified(partition.to_owned())
.await
.expect("recovered evidence verifies");
drop(host);
let _ = std::fs::remove_dir_all(dir);
}
#[test]
fn ambiguous_ensure_evicts_both_cached_handles_before_recovery() {
use commonware_runtime::{Runner as _, tokio as cw_tokio};
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-ensure-cache-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let storage_dir = dir.clone();
let partition = "conv-ensure-cache".to_owned();
let runner =
cw_tokio::Runner::new(cw_tokio::Config::default().with_storage_directory(dir.clone()));
runner.start(move |context| async move {
let mut logs = super::OpenLogs::new();
let mut checkpoints = HashMap::new();
let mut count_checkpoints = super::CountCheckpoints::new();
let mut mmr_logs = HashMap::new();
let signer = polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1);
let trust = polyc_crypto::signing_role::RoleTrustSet::from_public_keys(vec![
signer.public_key_bytes(),
])
.expect("test journal key is encoded ed25519");
arm_root_append_failure(&partition);
assert!(matches!(
super::append_batch_one(
&context,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&partition,
&[Event::new("output_msg", b"durable content".to_vec())],
)
.await,
Err(AppendError::AttestationOutcomeUnknown(_))
));
assert!(!logs.contains(&partition));
assert!(!mmr_logs.contains_key(&partition));
arm_root_append_failure(&partition);
assert!(matches!(
super::ensure_partition_attested_one(
&context,
&storage_dir,
&mut logs,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&trust,
&partition,
)
.await,
Err(AppendError::AttestationOutcomeUnknown(_))
));
assert!(
!logs.contains(&partition),
"an ambiguous ensure must not retain the possibly-fatal journal handle"
);
assert!(
!mmr_logs.contains_key(&partition),
"an ambiguous ensure must not retain the rebuilt MMR beside an uncertain append"
);
super::ensure_partition_attested_one(
&context,
&storage_dir,
&mut logs,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&trust,
&partition,
)
.await
.expect("a fresh reopen finishes attestation")
.expect("the durable content receives evidence");
let replayed = logs
.get_or_open(&context, &partition)
.await
.expect("reopen after recovery")
.replay()
.await
.expect("replay after recovery");
assert_eq!(
replayed
.iter()
.filter(|event| event.kind == polyc_eventlog::MMR_SIGNED_ROOT_KIND)
.count(),
1
);
});
let _ = std::fs::remove_dir_all(dir);
}
#[tokio::test]
async fn ensure_commit_ambiguity_settles_after_restart_without_a_second_root() {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-ensure-commit-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let partition = "conv-ensure-commit";
let signer = polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1);
let host = EventLogHost::spawn(dir.clone(), CancellationToken::new(), signer.clone())
.expect("spawn host");
arm_root_append_failure(partition);
assert!(matches!(
host.append_batch(
partition.to_owned(),
vec![Event::new("output_msg", b"durable content".to_vec())],
)
.await,
Err(AppendError::AttestationOutcomeUnknown(_))
));
arm_commit_failure(partition);
assert!(matches!(
host.ensure_partition_attested(partition.to_owned()).await,
Err(AppendError::AttestationOutcomeUnknown(_))
));
drop(host);
let reopened = EventLogHost::spawn(dir.clone(), CancellationToken::new(), signer)
.expect("restart host");
reopened
.ensure_partition_attested(partition.to_owned())
.await
.expect("same recovery settles after restart")
.expect("non-empty content is attested");
let replayed = reopened.replay(partition.to_owned()).await.expect("replay");
assert_eq!(
replayed
.iter()
.filter(|event| event.kind == polyc_eventlog::MMR_SIGNED_ROOT_KIND)
.count(),
1
);
reopened
.replay_verified(partition.to_owned())
.await
.expect("the sole recovered root verifies");
drop(reopened);
let _ = std::fs::remove_dir_all(dir);
}
#[derive(Clone, Copy, Debug)]
enum ExpectedOutcome {
AppendUnknown,
AttestationUnknown,
}
fn assert_expected_append_error(error: &AppendError, expected: ExpectedOutcome) {
match expected {
ExpectedOutcome::AppendUnknown => assert!(
matches!(error, AppendError::AppendOutcomeUnknown(_)),
"a failed event append cannot prove it left no recoverable bytes: {error}"
),
ExpectedOutcome::AttestationUnknown => assert!(
matches!(error, AppendError::AttestationOutcomeUnknown(_)),
"a final commit may already have made content and evidence durable: {error}"
),
}
}
fn poisoned_append_then_recovers(
dir_prefix: &str,
partition_name: &str,
failing_batch: Vec<Event>,
arm_failure: impl FnOnce(&str) + Send + 'static,
expected_outcome: ExpectedOutcome,
recovery_payload: &'static [u8],
) -> PoisonedAppendOutcome {
use commonware_runtime::{Runner as _, tokio as cw_tokio};
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-{dir_prefix}-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let _ = std::fs::remove_dir_all(&dir);
let partition = partition_name.to_owned();
let runner =
cw_tokio::Runner::new(cw_tokio::Config::default().with_storage_directory(dir.clone()));
let outcome = runner.start(move |context| async move {
let mut logs = super::OpenLogs::new();
let mut checkpoints = HashMap::new();
let mut count_checkpoints = super::CountCheckpoints::new();
let mut mmr_logs = HashMap::new();
let signer = polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1);
let trust = polyc_crypto::signing_role::RoleTrustSet::from_public_keys(vec![
signer.public_key_bytes(),
])
.expect("test journal key is encoded ed25519");
super::append_batch_one(
&context,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&partition,
&[Event::new("output_msg", b"seed".to_vec())],
)
.await
.expect("seed append");
assert!(
logs.contains(&partition),
"the seed append must have cached the handle"
);
let mmr_cached_after_seed = mmr_logs.contains_key(&partition);
arm_failure(&partition);
let err = super::append_batch_one(
&context,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&partition,
&failing_batch,
)
.await
.expect_err("the armed failure must surface");
assert_expected_append_error(&err, expected_outcome);
assert!(
!logs.contains(&partition),
"an injected failure must evict the partition's cached handle (#1712)"
);
let mmr_cached_after_failure = mmr_logs.contains_key(&partition);
super::append_batch_one(
&context,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&partition,
&[Event::new("output_msg", recovery_payload.to_vec())],
)
.await
.expect("a later append to the same partition must succeed cleanly");
let mmr_cached_after_recovery = mmr_logs.contains_key(&partition);
let log = logs
.get_or_open(&context, &partition)
.await
.expect("reopen");
let replayed = log.replay().await.expect("replay after recovery");
let replayed_payloads = replayed
.iter()
.map(|e| e.payload.clone())
.collect::<Vec<_>>();
let verified_after_recovery = super::replay_verified_one(
&context,
&mut logs,
&mut count_checkpoints,
&trust,
&partition,
)
.await;
PoisonedAppendOutcome {
replayed_payloads,
mmr_cached_after_seed,
mmr_cached_after_failure,
mmr_cached_after_recovery,
verified_after_recovery,
}
});
let _ = std::fs::remove_dir_all(&dir);
outcome
}
struct PoisonedAppendOutcome {
replayed_payloads: Vec<Vec<u8>>,
mmr_cached_after_seed: bool,
mmr_cached_after_failure: bool,
mmr_cached_after_recovery: bool,
verified_after_recovery: Result<Vec<Event>, super::VerifyError>,
}
#[test]
fn append_error_evicts_the_partitions_cached_handle() {
let outcome = poisoned_append_then_recovers(
"poison-append",
"conv-poison-append",
vec![
Event::new("output_msg", b"a".to_vec()),
Event::new("output_msg", b"b".to_vec()),
],
|partition| arm_append_failure(partition, 1),
ExpectedOutcome::AppendUnknown,
b"c",
);
let payloads: Vec<&[u8]> = outcome
.replayed_payloads
.iter()
.map(Vec::as_slice)
.collect();
assert!(
payloads.contains(&b"seed".as_slice()),
"the pre-existing seed event must survive"
);
assert!(
payloads.contains(&b"a".as_slice()),
"the first event of the failed batch was committed best-effort \
before the injected failure"
);
assert!(
!payloads.contains(&b"b".as_slice()),
"the event AT the injected failure must never have been appended"
);
assert!(
payloads.contains(&b"c".as_slice()),
"the recovery append must round-trip"
);
}
#[tokio::test]
async fn a_failed_commit_counts_no_attestation() {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-commit-fail-attest-{}",
std::process::id()
));
let _ = std::fs::remove_dir_all(&dir);
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(
dir.clone(),
shutdown.clone(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
)
.expect("spawn host");
host.append_batch(
"conv-commit-fail-attest".to_owned(),
vec![Event::new("output_msg", b"seed".to_vec())],
)
.await
.expect("seed");
let before = crate::metrics::attestation_count(crate::metrics::AttestationOutcome::Signed);
arm_commit_failure("conv-commit-fail-attest");
host.append_batch(
"conv-commit-fail-attest".to_owned(),
vec![Event::new("output_msg", b"never-lands".to_vec())],
)
.await
.expect_err("the armed commit failure refuses the batch");
assert_eq!(
crate::metrics::attestation_count(crate::metrics::AttestationOutcome::Signed),
before,
"a batch whose commit failed must not count as attested"
);
shutdown.cancel();
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn a_batch_that_fails_on_its_first_event_counts_nothing() {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-first-event-fail-{}",
std::process::id()
));
let _ = std::fs::remove_dir_all(&dir);
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(
dir.clone(),
shutdown.clone(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
)
.expect("spawn host");
host.append_batch(
"conv-first-event-fail".to_owned(),
vec![Event::new("output_msg", b"seed".to_vec())],
)
.await
.expect("seed");
let before = crate::metrics::attestation_count(
crate::metrics::AttestationOutcome::UnattestedPartialBatch,
);
arm_append_failure("conv-first-event-fail", 0);
host.append_batch(
"conv-first-event-fail".to_owned(),
vec![
Event::new("output_msg", b"never".to_vec()),
Event::new("output_msg", b"lands".to_vec()),
],
)
.await
.expect_err("the armed failure refuses the batch");
assert_eq!(
crate::metrics::attestation_count(
crate::metrics::AttestationOutcome::UnattestedPartialBatch
),
before,
"a batch where nothing landed leaves nothing unattested"
);
shutdown.cancel();
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn a_committed_partial_batch_counts_as_unattested() {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-partial-unattested-{}",
std::process::id()
));
let _ = std::fs::remove_dir_all(&dir);
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(
dir.clone(),
shutdown.clone(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
)
.expect("spawn host");
host.append_batch(
"conv-partial-unattested".to_owned(),
vec![Event::new("output_msg", b"seed".to_vec())],
)
.await
.expect("seed");
let before = crate::metrics::attestation_count(
crate::metrics::AttestationOutcome::UnattestedPartialBatch,
);
arm_append_failure("conv-partial-unattested", 1);
host.append_batch(
"conv-partial-unattested".to_owned(),
vec![
Event::new("output_msg", b"lands".to_vec()),
Event::new("output_msg", b"fails".to_vec()),
],
)
.await
.expect_err("the armed failure refuses the batch");
assert_eq!(
crate::metrics::attestation_count(
crate::metrics::AttestationOutcome::UnattestedPartialBatch
),
before + 1,
"a durable partial batch that never reached signing must be counted"
);
shutdown.cancel();
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn batch_order_first_element_survives_a_failure_on_the_second() {
let outcome = poisoned_append_then_recovers(
"poison-batch-order",
"conv-poison-batch-order",
vec![
Event::new("summary_gate_admitted", b"admission".to_vec()),
Event::new("summary", b"summary".to_vec()),
],
|partition| arm_append_failure(partition, 1),
ExpectedOutcome::AppendUnknown,
b"next-turn",
);
let payloads: Vec<&[u8]> = outcome
.replayed_payloads
.iter()
.map(Vec::as_slice)
.collect();
assert!(
payloads.contains(&b"admission".as_slice()),
"the admission event (ordered first) must survive the partial write"
);
assert!(
!payloads.contains(&b"summary".as_slice()),
"the summary event (ordered last) must never land without its \
paired admission — a partial write here must never leave a \
summary durable with no admission event anywhere in the log"
);
assert!(
payloads.contains(&b"next-turn".as_slice()),
"the partition recovers cleanly on the next append"
);
}
#[test]
fn commit_error_evicts_the_partitions_cached_handle() {
let outcome = poisoned_append_then_recovers(
"poison-commit",
"conv-poison-commit",
vec![Event::new("output_msg", b"never-committed".to_vec())],
arm_commit_failure,
ExpectedOutcome::AttestationUnknown,
b"after-recovery",
);
let payloads: Vec<&[u8]> = outcome
.replayed_payloads
.iter()
.map(Vec::as_slice)
.collect();
assert!(
payloads.contains(&b"seed".as_slice()),
"the pre-existing seed event must survive"
);
assert!(
payloads.contains(&b"after-recovery".as_slice()),
"the recovery append must round-trip"
);
}
#[test]
fn commit_failure_evicts_the_mmr_cache() {
let outcome = poisoned_append_then_recovers(
"poison-commit-mmr",
"conv-poison-commit-mmr",
vec![Event::new("output_msg", b"never-committed".to_vec())],
arm_commit_failure,
ExpectedOutcome::AttestationUnknown,
b"after-recovery",
);
assert!(
outcome.mmr_cached_after_seed,
"seed append must populate the running MMR"
);
assert!(
!outcome.mmr_cached_after_failure,
"#1740: a failed final commit must evict the partition's cached MMR"
);
assert!(
outcome.mmr_cached_after_recovery,
"#1740: a subsequent successful append must repopulate the MMR cache"
);
let verified = outcome.verified_after_recovery.expect(
"#1740: the partition's full history after a commit-failure-triggered \
MMR rebuild must still pass real tamper-evidence verification — a \
signed root computed post-recovery must check out against what \
actually landed on disk, not a phantom pre-failure tree",
);
assert!(
verified
.iter()
.any(|event| event.kind == polyc_eventlog::MMR_SIGNED_ROOT_KIND),
"the recovery append must have landed its own signed-root marker \
for the verification above to be meaningful"
);
}
#[test]
fn append_error_evicts_the_mmr_cache() {
let outcome = poisoned_append_then_recovers(
"poison-append-mmr",
"conv-poison-append-mmr",
vec![
Event::new("output_msg", b"a".to_vec()),
Event::new("output_msg", b"b".to_vec()),
],
|partition| arm_append_failure(partition, 1),
ExpectedOutcome::AppendUnknown,
b"after-recovery",
);
assert!(
outcome.mmr_cached_after_seed,
"seed append must populate the running MMR"
);
assert!(
!outcome.mmr_cached_after_failure,
"#1788: a failed per-event append must evict the partition's cached MMR"
);
assert!(
outcome.mmr_cached_after_recovery,
"#1788: a subsequent successful append must repopulate the MMR cache"
);
let verified = outcome.verified_after_recovery.expect(
"#1788: the partition's full history after a per-event-append-failure-triggered \
MMR rebuild must still pass real tamper-evidence verification — a \
signed root computed post-recovery must check out against what \
actually landed on disk, not a phantom pre-failure tree",
);
assert!(
verified
.iter()
.any(|event| event.kind == polyc_eventlog::MMR_SIGNED_ROOT_KIND),
"the recovery append must have landed its own signed-root marker \
for the verification above to be meaningful"
);
}
#[test]
fn destroy_failure_leaves_caches_intact() {
use commonware_runtime::{Runner as _, tokio as cw_tokio};
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-poison-destroy-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let _ = std::fs::remove_dir_all(&dir);
let partition = "conv-poison-destroy".to_owned();
let runner =
cw_tokio::Runner::new(cw_tokio::Config::default().with_storage_directory(dir.clone()));
runner.start(move |context| async move {
let mut logs = super::OpenLogs::new();
let mut checkpoints = HashMap::new();
let mut count_checkpoints = super::CountCheckpoints::new();
let mut mmr_logs = HashMap::new();
let signer = polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1);
super::append_batch_one(
&context,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&partition,
&[Event::new("output_msg", b"seed".to_vec())],
)
.await
.expect("seed append");
assert!(
mmr_logs.contains_key(&partition),
"seed append must populate the running MMR"
);
assert!(
count_checkpoints.contains(&partition),
"seed append must populate the durable count checkpoint"
);
checkpoints.insert(
partition.clone(),
CheckpointFrame {
tail_offset: 1,
payload: b"checkpoint".to_vec(),
},
);
arm_destroy_failure(&partition);
let err = super::destroy_one(
&context,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&partition,
super::DestroyCheckpoint::Clear,
)
.await
.expect_err("the armed destroy failure must surface");
assert!(
matches!(err, EventLogError::Journal(_)),
"expected a Journal error, got: {err}"
);
assert!(
checkpoints.contains_key(&partition),
"#1742: a failed destroy must not clear the compaction checkpoint frame"
);
assert!(
mmr_logs.contains_key(&partition),
"#1742: a failed destroy must not clear the running MMR"
);
assert!(
count_checkpoints.contains(&partition),
"#1742: a failed destroy must not clear the durable event-count checkpoint"
);
});
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn migrate_append_failure_evicts_the_handle_and_commits_partial_writes() {
use commonware_runtime::{Runner as _, tokio as cw_tokio};
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-poison-migrate-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let _ = std::fs::remove_dir_all(&dir);
let to = "conv-poison-migrate-to".to_owned();
let from = "conv-poison-migrate-from".to_owned();
let to_ref = super::PartitionRef::new(&to).expect("the name encodes");
let from_ref = super::PartitionRef::new(&from).expect("the name encodes");
let runner =
cw_tokio::Runner::new(cw_tokio::Config::default().with_storage_directory(dir.clone()));
runner.start(move |context| async move {
let mut logs = super::OpenLogs::new();
let mut count_checkpoints = super::CountCheckpoints::new();
let mut mmr_logs = HashMap::new();
let write_observers = super::observer::WriteObservers::new();
let mut checkpoints = HashMap::new();
super::append_batch_one(
&context,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
&to,
&[Event::new("output_msg", b"seed".to_vec())],
)
.await
.expect("seed the destination so its tree is cached");
assert!(
mmr_logs.contains_key(&to),
"the seed must leave a cached tree, or the assertion below is vacuous"
);
arm_append_failure(&to, 1);
let err = super::migrate_append_into(
&context,
&mut logs,
&mut count_checkpoints,
&mut mmr_logs,
&polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
&to_ref,
&[
Event::new("memory_added", b"a".to_vec()),
Event::new("memory_added", b"b".to_vec()),
],
&write_observers,
&from_ref,
)
.await
.expect_err("the armed failure must surface");
assert!(
matches!(err, EventLogError::Journal(_)),
"expected a Journal error, got: {err}"
);
assert!(
!logs.contains(&to),
"#1741: an injected append failure must evict `to`'s cached handle"
);
let recovered = logs.get_or_open(&context, &to).await.expect("reopen");
let replayed = recovered.replay().await.expect("replay");
let payloads: Vec<&[u8]> = replayed.iter().map(|e| e.payload.as_slice()).collect();
assert!(
payloads.contains(&b"a".as_slice()),
"the first event of the failed batch was committed best-effort"
);
assert!(
!payloads.contains(&b"b".as_slice()),
"the event AT the injected failure must never have been appended"
);
assert!(
!mmr_logs.contains_key(&to),
"an injected append failure must evict `to`'s cached MMR too"
);
});
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn migrate_commit_failure_evicts_the_handle() {
use commonware_runtime::{Runner as _, tokio as cw_tokio};
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-poison-migrate-commit-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let _ = std::fs::remove_dir_all(&dir);
let to = "conv-poison-migrate-commit-to".to_owned();
let from = "conv-poison-migrate-commit-from".to_owned();
let to_ref = super::PartitionRef::new(&to).expect("the name encodes");
let from_ref = super::PartitionRef::new(&from).expect("the name encodes");
let runner =
cw_tokio::Runner::new(cw_tokio::Config::default().with_storage_directory(dir.clone()));
runner.start(move |context| async move {
let mut logs = super::OpenLogs::new();
let mut count_checkpoints = super::CountCheckpoints::new();
let mut mmr_logs = HashMap::new();
let write_observers = super::observer::WriteObservers::new();
let mut checkpoints = HashMap::new();
super::append_batch_one(
&context,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
&to,
&[Event::new("output_msg", b"seed".to_vec())],
)
.await
.expect("seed the destination so its tree is cached");
assert!(
mmr_logs.contains_key(&to),
"the seed must leave a cached tree, or the assertion below is vacuous"
);
arm_commit_failure(&to);
let err = super::migrate_append_into(
&context,
&mut logs,
&mut count_checkpoints,
&mut mmr_logs,
&polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
&to_ref,
&[Event::new("memory_added", b"never-committed".to_vec())],
&write_observers,
&from_ref,
)
.await
.expect_err("the armed commit failure must surface");
assert!(
matches!(err, EventLogError::Journal(_)),
"expected a Journal error, got: {err}"
);
assert!(
!logs.contains(&to),
"#1741: an injected commit failure must evict `to`'s cached handle"
);
assert!(
!mmr_logs.contains_key(&to),
"an injected commit failure must evict `to`'s cached MMR too"
);
});
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn worker_label_interns_per_partition_and_is_reused_across_respawns() {
let mut workers = super::PartitionWorkers::new();
let a1 = workers.worker_label("conv-a");
let b1 = workers.worker_label("conv-b");
assert!(
!std::ptr::eq(a1, b1),
"distinct partitions must get distinct labels"
);
for _ in 0..50 {
assert!(
std::ptr::eq(a1, workers.worker_label("conv-a")),
"respawning `conv-a` must reuse its one interned label"
);
assert!(
std::ptr::eq(b1, workers.worker_label("conv-b")),
"respawning `conv-b` must reuse its one interned label"
);
}
assert_eq!(
workers.worker_labels.len(),
2,
"exactly one leaked label per DISTINCT partition, however many \
times each has been (re)spawned — not one per spawn"
);
}
#[test]
fn get_or_open_label_interns_per_partition_and_is_reused_across_reopens() {
use commonware_runtime::{Runner as _, tokio as cw_tokio};
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-label-intern-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let _ = std::fs::remove_dir_all(&dir);
let partition = "conv-label-intern".to_owned();
let runner =
cw_tokio::Runner::new(cw_tokio::Config::default().with_storage_directory(dir.clone()));
let distinct_labels = runner.start(move |context| async move {
let mut logs = super::OpenLogs::new();
logs.get_or_open(&context, &partition)
.await
.expect("first open");
let first_label = *logs
.log_labels
.get(&partition)
.expect("label recorded on first open");
for _ in 0..50 {
logs.take(&partition);
logs.get_or_open(&context, &partition)
.await
.expect("reopen after simulated eviction");
let label = *logs
.log_labels
.get(&partition)
.expect("label recorded on reopen");
assert!(
std::ptr::eq(first_label, label),
"reopening the SAME partition must reuse its one \
interned label"
);
}
logs.log_labels.len()
});
let _ = std::fs::remove_dir_all(&dir);
assert_eq!(
distinct_labels, 1,
"exactly one leaked label for the one distinct partition \
reopened, however many times it was evicted and reopened — \
not one per reopen"
);
}
#[tokio::test]
async fn append_rejects_payload_over_decode_cap() {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-cap-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(
dir.clone(),
shutdown,
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
)
.expect("spawn host");
let oversized = vec![0u8; polyc_proto::events_decode::MAX_EVENT_PAYLOAD_BYTES + 1];
let err = host
.append_batch(
"conv-cap-test".to_owned(),
vec![Event::new("output_msg", oversized)],
)
.await
.expect_err("over-cap payload must be refused");
assert!(
matches!(err, AppendError::PayloadTooLarge { ref kind, .. } if kind == "output_msg"),
"expected PayloadTooLarge, got: {err}"
);
drop(host);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn events_since_checkpoint_counts_only_the_post_checkpoint_tail() {
use polyc_proto::kinds::{
COMPACTION_CHECKPOINT as BASE_COMPACTION_CHECKPOINT, tagged as tagged_kind,
};
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-since-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(
dir.clone(),
shutdown,
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
)
.expect("spawn host");
let partition = "conv-since".to_owned();
host.append_batch(
partition.clone(),
vec![
Event::new("output_msg", b"a".to_vec()),
Event::new("output_msg", b"b".to_vec()),
Event::new("output_msg", b"c".to_vec()),
],
)
.await
.expect("append prefix");
assert_eq!(
host.events_since_checkpoint(partition.clone())
.await
.expect("count"),
4,
"no checkpoint ⇒ count is the whole log (3 content events + \
the batch's signed-root marker, #799)"
);
host.append_batch(
partition.clone(),
vec![Event::new(
tagged_kind(BASE_COMPACTION_CHECKPOINT, &uuid::Uuid::now_v7()),
b"ckpt".to_vec(),
)],
)
.await
.expect("append checkpoint");
host.append_batch(
partition.clone(),
vec![
Event::new("output_msg", b"d".to_vec()),
Event::new("output_msg", b"e".to_vec()),
],
)
.await
.expect("append tail");
let since = host
.events_since_checkpoint(partition.clone())
.await
.expect("count after checkpoint");
assert_eq!(since, 4, "count resets to the post-checkpoint tail, not 9");
let replay = host
.replay_tail_from_checkpoint(partition)
.await
.expect("replay tail");
assert_eq!(since, replay.tail.len() as u64);
drop(host);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn partition_event_count_tracks_total_appends_across_a_checkpoint() {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-count-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(
dir.clone(),
shutdown,
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
)
.expect("spawn host");
let partition = "conv-count".to_owned();
host.append_batch(
partition.clone(),
vec![
Event::new("output_msg", b"a".to_vec()),
Event::new("output_msg", b"b".to_vec()),
Event::new("output_msg", b"c".to_vec()),
],
)
.await
.expect("append prefix");
assert_eq!(
host.partition_event_count(partition.clone())
.await
.expect("count"),
4,
);
assert_eq!(
host.replay(partition.clone()).await.expect("replay").len() as u64,
4,
"partition_event_count must agree with a full replay's length"
);
host.append_batch(
partition.clone(),
vec![Event::new(
polyc_proto::kinds::tagged(
polyc_proto::kinds::COMPACTION_CHECKPOINT,
&uuid::Uuid::now_v7(),
),
b"ckpt".to_vec(),
)],
)
.await
.expect("append checkpoint");
host.append_batch(
partition.clone(),
vec![
Event::new("output_msg", b"d".to_vec()),
Event::new("output_msg", b"e".to_vec()),
],
)
.await
.expect("append tail");
let total = host
.partition_event_count(partition.clone())
.await
.expect("count after checkpoint");
assert_eq!(total, 9, "total keeps growing across a checkpoint");
let since = host
.events_since_checkpoint(partition.clone())
.await
.expect("since checkpoint, for contrast");
assert_ne!(
total, since,
"the two primitives measure different things once a checkpoint exists"
);
drop(host);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn partition_event_count_of_unknown_partition_is_zero_and_creates_nothing() {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-count-unknown-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let _ = std::fs::remove_dir_all(&dir);
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(
dir.clone(),
shutdown,
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
)
.expect("spawn host");
let count = host
.partition_event_count("conv-never-touched".to_owned())
.await
.expect("count of unknown partition");
assert_eq!(count, 0, "an unwritten partition counts as empty");
assert!(
!dir.join("conv-never-touched_data").exists(),
"a mere count probe must not create the partition's journal"
);
assert!(
!dir.join("conv-never-touched_offsets").exists(),
"a mere count probe must not create the partition's offsets index either"
);
drop(host);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn partition_last_modified_reports_storage_time_and_creates_nothing() {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-lastmod-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let _ = std::fs::remove_dir_all(&dir);
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(
dir.clone(),
shutdown,
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
)
.expect("spawn host");
host.append_batch(
"conv-app:one".to_owned(),
vec![Event::new("k", b"v".to_vec())],
)
.await
.expect("seed");
let written = host
.partition_last_modified_ms("conv-app:one".to_owned())
.await
.expect("last-modified of a written partition")
.expect("a written partition has a storage time");
let now = u64::try_from(
std::time::SystemTime::now()
.duration_since(std::time::SystemTime::UNIX_EPOCH)
.expect("clock")
.as_millis(),
)
.expect("a millisecond clock reading fits");
assert!(
now.saturating_sub(written) < 60_000,
"a partition written moments ago reports a recent time: {written} against {now}"
);
assert_eq!(
host.partition_last_modified_ms("conv-never-touched".to_owned())
.await
.expect("last-modified of an unknown partition"),
None,
"a partition this host holds nothing for reports nothing"
);
assert!(
!dir.join("conv-never-touched_data").exists(),
"a mere time probe must not create the partition's journal"
);
let long_ago = std::time::SystemTime::now() - std::time::Duration::from_hours(1);
let data_dir = dir.join(format!(
"{}_data",
super::encode_partition("conv-app:one").expect("the name encodes")
));
for entry in std::fs::read_dir(&data_dir).expect("read data dir") {
let entry = entry.expect("dir entry");
std::fs::File::options()
.write(true)
.open(entry.path())
.expect("open for mtime set")
.set_modified(long_ago)
.expect("set_modified");
}
let backdated = host
.partition_last_modified_ms("conv-app:one".to_owned())
.await
.expect("last-modified after backdating")
.expect("a backdated partition still has a storage time");
assert!(
backdated < written,
"the reported time follows the storage it is read from: {backdated} against {written}"
);
drop(host);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn mmr_append_batch_signs_root_and_replay_verified_survives_reopen() {
let dir =
std::env::temp_dir().join(format!("polychrome-eventlog-mmr-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(
dir.clone(),
shutdown.clone(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
)
.expect("spawn");
let partition = "conv-mmr-wiring".to_owned();
host.append_batch(
partition.clone(),
vec![Event::new("user_msg", b"turn one".to_vec())],
)
.await
.expect("append turn 1");
host.append_batch(
partition.clone(),
vec![Event::new("output_msg", b"turn two".to_vec())],
)
.await
.expect("append turn 2");
let verified = host
.replay_verified(partition.clone())
.await
.expect("both turns verify within the same process");
let marker_count = verified
.iter()
.filter(|e| e.kind == polyc_eventlog::MMR_SIGNED_ROOT_KIND)
.count();
assert_eq!(
verified.len(),
4,
"2 content events + 2 signed-root markers"
);
assert_eq!(marker_count, 2, "one marker per turn's append_batch");
drop(host);
let shutdown2 = CancellationToken::new();
let host2 = EventLogHost::spawn(
dir.clone(),
shutdown2,
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
)
.expect("respawn");
host2
.append_batch(
partition.clone(),
vec![Event::new("tool_call", b"turn three".to_vec())],
)
.await
.expect("append turn 3 after reopen");
let verified_after_reopen = host2
.replay_verified(partition)
.await
.expect("all three turns verify after a cold rebuild");
assert_eq!(
verified_after_reopen.len(),
6,
"3 content events + 3 signed-root markers, across the reopen"
);
drop(host2);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn repair_apply_refuses_a_partition_changed_after_inspection() {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-repair-fence-{}",
std::process::id()
));
let _ = std::fs::remove_dir_all(&dir);
let partition = "conv-repair-fence".to_owned();
let host = EventLogHost::spawn(
dir.clone(),
CancellationToken::new(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
)
.expect("spawn");
host.append_batch(
partition.clone(),
vec![Event::new("user_msg", b"before".to_vec())],
)
.await
.expect("seed");
let stale = host
.inspect_repair_partition(partition.clone(), None)
.await
.expect("inspect");
host.append_batch(
partition.clone(),
vec![Event::new("user_msg", b"after".to_vec())],
)
.await
.expect("concurrent append");
assert_eq!(
host.apply_repair_partition(partition.clone(), stale.before, stale.after, None)
.await
.expect("conditioned apply"),
RepairApplication::Changed
);
assert_eq!(
host.replay_verified(partition)
.await
.expect("changed partition remains intact")
.iter()
.filter(|event| event.kind == "user_msg")
.count(),
2
);
drop(host);
let _ = std::fs::remove_dir_all(dir);
}
#[allow(clippy::too_many_lines)]
#[test]
fn repair_resumes_from_staged_survivors_after_restore_fails() {
use commonware_runtime::{Runner as _, tokio as cw_tokio};
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-repair-stage-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let _ = std::fs::remove_dir_all(&dir);
let partition = "conv-repair-stage".to_owned();
let run_dir = dir.clone();
let runner =
cw_tokio::Runner::new(cw_tokio::Config::default().with_storage_directory(dir.clone()));
runner.start(move |context| async move {
let mut logs = super::OpenLogs::new();
let mut checkpoints = HashMap::new();
let mut count_checkpoints = super::CountCheckpoints::new();
let mut mmr_logs = HashMap::new();
let signer = polyc_crypto::signing_role::JournalAttestationSigner::from_seed(7);
let trust = polyc_crypto::signing_role::RoleTrustSet::current(&signer);
let original = [
Event::new("user_msg", b"survivor".to_vec()),
Event::new("broken", b"quarantine".to_vec()),
];
{
let log = logs.get_or_open(&context, &partition).await.unwrap();
for event in &original {
log.append(event).await.unwrap();
}
log.commit().await.unwrap();
}
count_checkpoints
.record(&context, &partition, 2)
.await
.unwrap();
let inspected = super::inspect_repair_one(
&context,
&mut logs,
&mut count_checkpoints,
&signer,
&trust,
&partition,
None,
)
.await
.unwrap();
let survivor = inspected.readable[0].clone();
let (restored_content, _) =
super::repaired_content(std::slice::from_ref(&survivor.1), &signer, None)
.expect("re-root");
let indexed: Vec<(u64, Event)> = restored_content
.iter()
.enumerate()
.map(|(position, event)| (position as u64, event.clone()))
.collect();
let after = super::repair_fingerprint(&indexed, &[], restored_content.len() as u64);
let scan = super::RepairScan {
readable: vec![survivor],
quarantined: vec![polyc_eventlog::QuarantinedItem {
position: 1,
error: "injected unreadable frame".to_owned(),
}],
before: inspected.before,
after,
standing: inspected.standing,
};
let stage = super::repair_stage_partition(&partition, &inspected.before);
arm_append_failure(&stage, 1);
super::apply_repair_scan(
&context,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&partition,
scan.clone(),
None,
)
.await
.expect_err("the armed staging failure surfaces");
assert_eq!(
logs.get_or_open(&context, &partition)
.await
.unwrap()
.replay()
.await
.unwrap(),
original,
"a stage without its sentinel cannot authorize target destruction"
);
arm_append_failure(&partition, 0);
super::apply_repair_scan(
&context,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&partition,
scan,
None,
)
.await
.expect_err("the armed restore failure surfaces");
assert!(
super::list_partitions_in(&run_dir)
.unwrap()
.iter()
.all(|name| !name.starts_with(super::REPAIR_STAGE_PREFIX)),
"repair staging is never a caller-visible partition"
);
logs = super::OpenLogs::new();
checkpoints.clear();
count_checkpoints.evict_all();
mmr_logs.clear();
let (application, _) = super::apply_repair_conditionally(
&context,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&trust,
&partition,
super::RepairFingerprints {
before: inspected.before,
after,
},
None,
)
.await
.expect("restart resumes from the durable stage");
assert_eq!(application, RepairApplication::Applied);
let restored = logs
.get_or_open(&context, &partition)
.await
.unwrap()
.replay()
.await
.unwrap();
assert_eq!(
restored, restored_content,
"a resumed repair restores the re-rooted content, marker included"
);
assert_eq!(
restored.len(),
2,
"one survivor plus exactly one fresh signed root"
);
assert_eq!(restored[0], original[0]);
assert_eq!(restored[1].kind, polyc_eventlog::MMR_SIGNED_ROOT_KIND);
assert!(
super::list_partitions_in(&run_dir)
.unwrap()
.iter()
.all(|name| !name.starts_with(super::REPAIR_STAGE_PREFIX)),
"the completed repair removes its private stage"
);
});
let _ = std::fs::remove_dir_all(dir);
}
#[test]
#[allow(clippy::too_many_lines)]
fn one_decision_drives_inspection_the_plan_and_application() {
use commonware_runtime::{Runner as _, tokio as cw_tokio};
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-one-decision-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let _ = std::fs::remove_dir_all(&dir);
let run_dir = dir.clone();
let runner =
cw_tokio::Runner::new(cw_tokio::Config::default().with_storage_directory(dir.clone()));
runner.start(move |context| async move {
let mut logs = super::OpenLogs::new();
let mut checkpoints = HashMap::new();
let mut count_checkpoints = super::CountCheckpoints::new();
let mut mmr_logs = HashMap::new();
let signer = polyc_crypto::signing_role::JournalAttestationSigner::from_seed(7);
let trust = polyc_crypto::signing_role::RoleTrustSet::current(&signer);
let observers = super::observer::WriteObservers::new();
let partition = "conv-one-decision";
let content = [
Event::new("journal_lineage", b"old-incarnation".to_vec()),
Event::new("user_msg", b"one".to_vec()),
Event::new("output_msg", b"two".to_vec()),
];
let marker = polyc_eventlog::extend_and_sign(
&polyc_mmr::VerifiableLog::new(),
&content,
&signer,
)
.expect("sign over three leaves");
{
let log = logs.get_or_open(&context, partition).await.unwrap();
for event in [&content[0], &content[1], &marker] {
log.append(event).await.unwrap();
}
log.commit().await.unwrap();
}
count_checkpoints
.record(&context, partition, 3)
.await
.unwrap();
let replacement =
super::RepairEventReplacement::new("journal_lineage", b"new-incarnation".to_vec());
let (ack, ack_rx) = tokio::sync::oneshot::channel();
super::handle_command(
&context,
&run_dir,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&trust,
&[],
&observers,
super::Command::InspectRepair {
partition: super::PartitionRef::new(partition).expect("a storable name"),
replacement: Some(replacement.clone()),
ack,
},
)
.await;
let plan = ack_rx.await.expect("the arm acks").expect("inspectable");
assert!(
plan.quarantined.is_empty(),
"the premise is a stale root with nothing unreadable: {:?}",
plan.quarantined
);
assert!(
plan.will_rewrite,
"a stale root re-roots, so the durable plan must record a rewrite even with an \
empty quarantine list"
);
let (ack, ack_rx) = tokio::sync::oneshot::channel();
super::handle_command(
&context,
&run_dir,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&trust,
&[],
&observers,
super::Command::ApplyRepair {
partition: super::PartitionRef::new(partition).expect("a storable name"),
before: plan.before,
after: plan.after,
replacement: Some(replacement.clone()),
ack,
},
)
.await;
assert_eq!(
ack_rx.await.expect("the arm acks").expect("appliable"),
super::RepairApplication::Applied,
"application must produce exactly the `after` inspection promised. `Changed` here \
means the two phases decided the rewrite differently"
);
let restored = super::replay_verified_one(
&context,
&mut logs,
&mut count_checkpoints,
&trust,
partition,
)
.await
.expect("the repaired partition verifies");
let lineages: Vec<Vec<u8>> = restored
.iter()
.filter(|event| event.kind == "journal_lineage")
.map(|event| event.payload.clone())
.collect();
assert_eq!(
lineages,
vec![b"new-incarnation".to_vec()],
"the replacement applies exactly once, so the partition carries the rotated \
lineage and not the old one"
);
let (ack, ack_rx) = tokio::sync::oneshot::channel();
super::handle_command(
&context,
&run_dir,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&trust,
&[],
&observers,
super::Command::RepairSettled {
partition: super::PartitionRef::new(partition).expect("a storable name"),
before: plan.before,
after: plan.after,
ack,
},
)
.await;
assert!(
ack_rx.await.expect("the arm acks").expect("settled reads"),
"the destructive half is complete, so the repair reads as settled"
);
});
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
#[allow(clippy::too_many_lines)]
fn repair_re_roots_a_stale_root_leaves_a_healthy_partition_alone_and_refuses_a_tampered_one() {
use commonware_runtime::{Runner as _, tokio as cw_tokio};
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-stale-root-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let _ = std::fs::remove_dir_all(&dir);
let run_dir = dir.clone();
let runner =
cw_tokio::Runner::new(cw_tokio::Config::default().with_storage_directory(dir.clone()));
runner.start(move |context| async move {
let mut logs = super::OpenLogs::new();
let mut checkpoints = HashMap::new();
let mut count_checkpoints = super::CountCheckpoints::new();
let mut mmr_logs = HashMap::new();
let signer = polyc_crypto::signing_role::JournalAttestationSigner::from_seed(3);
let trust = polyc_crypto::signing_role::RoleTrustSet::current(&signer);
let content = [
Event::new("user_msg", b"one".to_vec()),
Event::new("output_msg", b"two".to_vec()),
Event::new("output_msg", b"three".to_vec()),
];
let marker = polyc_eventlog::extend_and_sign(
&polyc_mmr::VerifiableLog::new(),
&content,
&signer,
)
.expect("sign over three leaves");
let seed = async |logs: &mut super::OpenLogs<_>,
count_checkpoints: &mut super::CountCheckpoints<_>,
partition: &str,
events: &[Event]| {
{
let log = logs.get_or_open(&context, partition).await.unwrap();
for event in events {
log.append(event).await.unwrap();
}
log.commit().await.unwrap();
}
count_checkpoints
.record(&context, partition, events.len() as u64)
.await
.unwrap();
};
let stale = "conv-stale-root";
seed(
&mut logs,
&mut count_checkpoints,
stale,
&[content[0].clone(), content[1].clone(), marker.clone()],
)
.await;
super::replay_verified_one(&context, &mut logs, &mut count_checkpoints, &trust, stale)
.await
.expect_err("the seeded partition does not verify — that is the premise");
let observers = super::observer::WriteObservers::new();
let recorder =
std::sync::Arc::new(super::observer::tests::RecordingObserver::default());
observers.register(recorder.clone());
assert_eq!(observers.epoch(stale), 0);
let (ack, ack_rx) = tokio::sync::oneshot::channel();
super::handle_command(
&context,
&run_dir,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&trust,
&[],
&observers,
super::Command::RepairPartition {
partition: super::PartitionRef::new(stale).expect("a storable name"),
ack,
},
)
.await;
let quarantined = ack_rx.await.expect("the arm acks").expect("repairable");
assert!(
quarantined.is_empty(),
"nothing was unreadable, so nothing is dropped: {quarantined:?}"
);
assert_eq!(
observers.epoch(stale),
1,
"a repair that only re-roots still changed the partition, so every cache keyed \
on this epoch must be invalidated"
);
let notified = recorder
.mutations
.lock()
.expect("poison")
.iter()
.map(|mutation| mutation.1.clone())
.collect::<Vec<_>>();
assert_eq!(
notified,
vec![super::MutationKind::Repaired { quarantined: 0 }],
"one notification for the re-root, and it reports nothing dropped"
);
let repaired = super::replay_verified_one(
&context,
&mut logs,
&mut count_checkpoints,
&trust,
stale,
)
.await
.expect("the repaired partition verifies");
assert_eq!(
repaired
.iter()
.filter(|event| event.kind != polyc_eventlog::MMR_SIGNED_ROOT_KIND)
.count(),
2,
"re-rooting keeps every readable event"
);
let healthy = "conv-healthy";
let first = [Event::new("user_msg", b"a".to_vec())];
let first_marker = super::append_batch_one(
&context,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
healthy,
&first,
)
.await
.expect("first batch");
let _ = first_marker;
super::append_batch_one(
&context,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
healthy,
&[Event::new("output_msg", b"b".to_vec())],
)
.await
.expect("second batch");
let before = logs
.get_or_open(&context, healthy)
.await
.unwrap()
.replay()
.await
.unwrap();
assert_eq!(
before
.iter()
.filter(|event| event.kind == polyc_eventlog::MMR_SIGNED_ROOT_KIND)
.count(),
2,
"two batches, two roots — the state a naive trigger would rewrite"
);
let (healthy_ack, healthy_rx) = tokio::sync::oneshot::channel();
super::handle_command(
&context,
&run_dir,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&trust,
&[],
&observers,
super::Command::RepairPartition {
partition: super::PartitionRef::new(healthy).expect("a storable name"),
ack: healthy_ack,
},
)
.await;
healthy_rx
.await
.expect("the arm acks")
.expect("a healthy partition repairs to a no-op");
assert_eq!(
observers.epoch(healthy),
0,
"a repair that changed nothing must not invalidate a single cache"
);
assert_eq!(
logs.get_or_open(&context, healthy)
.await
.unwrap()
.replay()
.await
.unwrap(),
before,
"a healthy partition is left byte-identical"
);
let tampered = "conv-tampered";
let mut altered = content[0].clone();
altered.payload[0] ^= 0xFF;
let pair = [content[0].clone(), content[1].clone()];
let pair_marker =
polyc_eventlog::extend_and_sign(&polyc_mmr::VerifiableLog::new(), &pair, &signer)
.expect("sign over two leaves");
seed(
&mut logs,
&mut count_checkpoints,
tampered,
&[altered, content[1].clone(), pair_marker],
)
.await;
super::repair_one(
&context,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&trust,
tampered,
)
.await
.expect_err("repair must not re-sign content a trusted root contradicts");
});
let _ = std::fs::remove_dir_all(dir);
}
#[test]
fn a_destroy_that_fails_still_records_that_it_started() {
use commonware_runtime::{Runner as _, tokio as cw_tokio};
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-destroy-order-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let _ = std::fs::remove_dir_all(&dir);
let runner =
cw_tokio::Runner::new(cw_tokio::Config::default().with_storage_directory(dir.clone()));
runner.start(move |context| async move {
let mut logs = super::OpenLogs::new();
let mut checkpoints = HashMap::new();
let mut count_checkpoints = super::CountCheckpoints::new();
let mut mmr_logs = HashMap::new();
let partition = "conv-destroy-order";
{
let log = logs.get_or_open(&context, partition).await.unwrap();
log.append(&Event::new("user_msg", b"content".to_vec()))
.await
.unwrap();
log.commit().await.unwrap();
}
count_checkpoints
.record(&context, partition, 1)
.await
.unwrap();
super::tests::arm_destroy_failure(partition);
let outcome = super::destroy_one(
&context,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
partition,
super::DestroyCheckpoint::Clear,
)
.await;
assert!(outcome.is_err(), "the armed destroy failure surfaces");
assert_eq!(
count_checkpoints
.expected_count(&context, partition)
.await
.unwrap(),
0,
"the floor is cleared before anything is removed, so a destroy that failed \
part-way still says it started"
);
assert!(
super::repair_window_is_open(&context, &mut count_checkpoints, partition)
.await
.unwrap(),
"and a settlement probe reads that record rather than guessing from content"
);
});
let _ = std::fs::remove_dir_all(dir);
}
#[test]
fn only_the_recorded_floor_decides_whether_a_repair_is_still_in_its_window() {
use commonware_runtime::{Runner as _, tokio as cw_tokio};
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-window-floor-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let _ = std::fs::remove_dir_all(&dir);
let runner =
cw_tokio::Runner::new(cw_tokio::Config::default().with_storage_directory(dir.clone()));
runner.start(move |context| async move {
let mut count_checkpoints = super::CountCheckpoints::new();
let partition = "conv-window-floor";
count_checkpoints
.record(&context, partition, 7)
.await
.unwrap();
assert!(
!super::repair_window_is_open(&context, &mut count_checkpoints, partition)
.await
.unwrap(),
"a recorded floor is a record that content is present"
);
count_checkpoints.clear(&context, partition).await.unwrap();
assert!(
super::repair_window_is_open(&context, &mut count_checkpoints, partition)
.await
.unwrap(),
"a cleared floor is the destroy's own record that content was taken away"
);
count_checkpoints
.record(&context, partition, 2)
.await
.unwrap();
assert!(
!super::repair_window_is_open(&context, &mut count_checkpoints, partition)
.await
.unwrap(),
"the terminal record closes the window"
);
});
let _ = std::fs::remove_dir_all(dir);
}
#[test]
#[allow(clippy::too_many_lines)]
fn the_erasure_sweep_removes_a_stage_and_is_safe_to_repeat() {
use commonware_runtime::{Runner as _, tokio as cw_tokio};
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-sweep-stage-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let _ = std::fs::remove_dir_all(&dir);
let partition = "conv-sweep-stage".to_owned();
let run_dir = dir.clone();
let runner =
cw_tokio::Runner::new(cw_tokio::Config::default().with_storage_directory(dir.clone()));
runner.start(move |context| async move {
let mut logs = super::OpenLogs::new();
let mut checkpoints = HashMap::new();
let mut count_checkpoints = super::CountCheckpoints::new();
let mut mmr_logs = HashMap::new();
let signer = polyc_crypto::signing_role::JournalAttestationSigner::from_seed(31);
let trust = polyc_crypto::signing_role::RoleTrustSet::current(&signer);
let observers = super::observer::WriteObservers::new();
let before = [9_u8; 32];
let after = [11_u8; 32];
let stage = super::repair_stage_partition(&partition, &before);
let other = super::repair_stage_partition(&partition, &[7_u8; 32]);
let staged = vec![
Event::new("user_msg", b"private".to_vec()),
super::repair_stage_complete(before, after, 1),
];
super::append_repair_copy(&context, &mut logs, &stage, &staged)
.await
.expect("plant the stage");
super::append_repair_copy(&context, &mut logs, &other, &staged)
.await
.expect("plant a second repair's stage");
assert!(
super::partition_data_dir_exists(&run_dir, &stage).unwrap(),
"the planted stage is on the volume for this test to mean anything"
);
macro_rules! sweep {
() => {{
let (ack, ack_rx) = tokio::sync::oneshot::channel();
super::handle_command(
&context,
&run_dir,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&trust,
&[],
&observers,
super::Command::DestroyRepairStage {
partition: super::PartitionRef::new(&partition)
.expect("a storable name"),
before,
ack,
},
)
.await;
ack_rx.await.expect("the arm acks")
}};
}
sweep!().expect("the sweep runs");
assert!(
!super::partition_data_dir_exists(&run_dir, &stage).unwrap(),
"the erased conversation's staged copy is gone from the volume"
);
sweep!().expect("a repeated sweep is not an error");
assert!(
!super::partition_data_dir_exists(&run_dir, &stage).unwrap(),
"and it does not create the stage it was asked to remove"
);
assert!(
super::partition_data_dir_exists(&run_dir, &other).unwrap(),
"a sweep removes the stage it names and no other"
);
let rewrite = super::rewrite_stage_partition(&partition, "cmd-excise");
super::append_repair_copy(&context, &mut logs, &rewrite, &staged)
.await
.expect("plant an excision's stage");
let (ack, ack_rx) = tokio::sync::oneshot::channel();
super::handle_command(
&context,
&run_dir,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&trust,
&[],
&observers,
super::Command::DestroyRewriteStage {
partition: super::PartitionRef::new(&partition).expect("a storable name"),
command_id: "cmd-excise".to_owned(),
ack,
},
)
.await;
ack_rx
.await
.expect("the arm acks")
.expect("the rewrite sweep runs");
assert!(
!super::partition_data_dir_exists(&run_dir, &rewrite).unwrap(),
"an excision's staged copy is gone from the volume too"
);
});
let _ = std::fs::remove_dir_all(dir);
}
#[test]
#[allow(clippy::too_many_lines)]
fn a_stage_that_outlived_its_repair_neither_refuses_nor_rolls_back_the_partition() {
use commonware_runtime::{Runner as _, tokio as cw_tokio};
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-stale-stage-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let _ = std::fs::remove_dir_all(&dir);
let partition = "conv-stale-stage".to_owned();
let run_dir = dir.clone();
let runner =
cw_tokio::Runner::new(cw_tokio::Config::default().with_storage_directory(dir.clone()));
runner.start(move |context| async move {
let mut logs = super::OpenLogs::new();
let mut checkpoints = HashMap::new();
let mut count_checkpoints = super::CountCheckpoints::new();
let mut mmr_logs = HashMap::new();
let signer = polyc_crypto::signing_role::JournalAttestationSigner::from_seed(17);
let trust = polyc_crypto::signing_role::RoleTrustSet::current(&signer);
let original = [
Event::new("user_msg", b"survivor".to_vec()),
Event::new("broken", b"quarantine".to_vec()),
];
{
let log = logs.get_or_open(&context, &partition).await.unwrap();
for event in &original {
log.append(event).await.unwrap();
}
log.commit().await.unwrap();
}
count_checkpoints
.record(&context, &partition, 2)
.await
.unwrap();
let inspected = super::inspect_repair_one(
&context,
&mut logs,
&mut count_checkpoints,
&signer,
&trust,
&partition,
None,
)
.await
.unwrap();
let survivor = inspected.readable[0].clone();
let (restored_content, _) =
super::repaired_content(std::slice::from_ref(&survivor.1), &signer, None)
.expect("re-root");
let indexed: Vec<(u64, Event)> = restored_content
.iter()
.enumerate()
.map(|(position, event)| (position as u64, event.clone()))
.collect();
let after = super::repair_fingerprint(&indexed, &[], restored_content.len() as u64);
let fingerprints = super::RepairFingerprints {
before: inspected.before,
after,
};
super::apply_repair_scan(
&context,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&partition,
super::RepairScan {
readable: vec![survivor],
quarantined: vec![polyc_eventlog::QuarantinedItem {
position: 1,
error: "injected unreadable frame".to_owned(),
}],
before: inspected.before,
after,
standing: inspected.standing,
},
None,
)
.await
.expect("the repair completes");
let stage = super::repair_stage_partition(&partition, &inspected.before);
let mut staged = restored_content.clone();
staged.push(super::repair_stage_complete(
inspected.before,
after,
restored_content.len() as u64,
));
super::append_repair_copy(&context, &mut logs, &stage, &staged)
.await
.expect("the stage is put back exactly as the repair wrote it");
assert!(
super::repair_settled_one(
&context,
&run_dir,
&mut logs,
&mut count_checkpoints,
&signer,
&trust,
&partition,
fingerprints,
)
.await
.expect("probe"),
"a repaired partition is settled even with its stage still on disk"
);
super::append_batch_one(
&context,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&partition,
&[Event::new("user_msg", b"after the repair".to_vec())],
)
.await
.expect("a repaired partition takes new commits");
let grown = logs
.get_or_open(&context, &partition)
.await
.unwrap()
.replay()
.await
.unwrap();
assert!(
super::repair_settled_one(
&context,
&run_dir,
&mut logs,
&mut count_checkpoints,
&signer,
&trust,
&partition,
fingerprints,
)
.await
.expect("probe"),
"a partition that grew past its repair is settled, stale stage or not — \
refusing it would take a healthy conversation offline for good"
);
let (application, _) = super::apply_repair_conditionally(
&context,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&trust,
&partition,
fingerprints,
None,
)
.await
.expect("the retry answers rather than failing");
assert_eq!(
application,
RepairApplication::Changed,
"the partition moved on, so the retry reports that rather than resuming"
);
assert_eq!(
logs.get_or_open(&context, &partition)
.await
.unwrap()
.replay()
.await
.unwrap(),
grown,
"every committed event survives the retry"
);
count_checkpoints
.record(&context, &partition, grown.len() as u64 + 2)
.await
.unwrap();
let second = super::repair_one(
&context,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&trust,
&partition,
)
.await
.expect("a second repair runs on a healthy partition");
assert!(
second.rewrote,
"this block is only meaningful if the second repair actually rewrote the \
partition: {second:?}"
);
let rewritten = logs
.get_or_open(&context, &partition)
.await
.unwrap()
.replay()
.await
.unwrap();
assert!(
super::repair_settled_one(
&context,
&run_dir,
&mut logs,
&mut count_checkpoints,
&signer,
&trust,
&partition,
fingerprints,
)
.await
.expect("probe"),
"a partition a later repair rewrote is settled — the first repair's stage \
describes none of it"
);
let (application, _) = super::apply_repair_conditionally(
&context,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&trust,
&partition,
fingerprints,
None,
)
.await
.expect("the retry answers rather than failing");
assert_eq!(
application,
RepairApplication::Changed,
"the first repair's stage must not be restored over a later repair's work"
);
assert_eq!(
logs.get_or_open(&context, &partition)
.await
.unwrap()
.replay()
.await
.unwrap(),
rewritten,
"the later repair's content survives the stale stage's retry"
);
});
let _ = std::fs::remove_dir_all(dir);
}
#[test]
#[allow(clippy::too_many_lines)]
fn a_repair_that_crossed_its_destructive_boundary_reads_as_unsettled() {
use commonware_runtime::{Runner as _, tokio as cw_tokio};
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-repair-settled-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let _ = std::fs::remove_dir_all(&dir);
let partition = "conv-repair-settled".to_owned();
let run_dir = dir.clone();
let runner =
cw_tokio::Runner::new(cw_tokio::Config::default().with_storage_directory(dir.clone()));
runner.start(move |context| async move {
let mut logs = super::OpenLogs::new();
let mut checkpoints = HashMap::new();
let mut count_checkpoints = super::CountCheckpoints::new();
let mut mmr_logs = HashMap::new();
let signer = polyc_crypto::signing_role::JournalAttestationSigner::from_seed(11);
let trust = polyc_crypto::signing_role::RoleTrustSet::current(&signer);
let original = [
Event::new("user_msg", b"survivor".to_vec()),
Event::new("broken", b"quarantine".to_vec()),
];
{
let log = logs.get_or_open(&context, &partition).await.unwrap();
for event in &original {
log.append(event).await.unwrap();
}
log.commit().await.unwrap();
}
count_checkpoints
.record(&context, &partition, 2)
.await
.unwrap();
let inspected = super::inspect_repair_one(
&context,
&mut logs,
&mut count_checkpoints,
&signer,
&trust,
&partition,
None,
)
.await
.unwrap();
let survivor = inspected.readable[0].clone();
let (restored_content, _) =
super::repaired_content(std::slice::from_ref(&survivor.1), &signer, None)
.expect("re-root");
let indexed: Vec<(u64, Event)> = restored_content
.iter()
.enumerate()
.map(|(position, event)| (position as u64, event.clone()))
.collect();
let after = super::repair_fingerprint(&indexed, &[], restored_content.len() as u64);
let fingerprints = super::RepairFingerprints {
before: inspected.before,
after,
};
macro_rules! settled {
() => {
super::repair_settled_one(
&context,
&run_dir,
&mut logs,
&mut count_checkpoints,
&signer,
&trust,
&partition,
fingerprints,
)
.await
.expect("the probe reads the stage")
};
}
assert!(
settled!(),
"a repair that never staged anything leaves nothing to distrust"
);
{
let scan = super::inspect_repair_one(
&context,
&mut logs,
&mut count_checkpoints,
&signer,
&trust,
&partition,
None,
)
.await
.unwrap();
assert_eq!(scan.before, fingerprints.before, "the target is untouched");
let (restored, _) =
super::repaired_content(std::slice::from_ref(&original[0]), &signer, None)
.unwrap();
let mut planted = restored.clone();
planted.push(super::repair_stage_complete(
scan.before,
after,
restored.len() as u64,
));
let stage = super::repair_stage_partition(&partition, &fingerprints.before);
super::append_repair_copy(&context, &mut logs, &stage, &planted)
.await
.expect("plant a complete stage beside an untouched target");
}
assert!(
settled!(),
"a complete stage over an untouched target is not a window: the partition is \
byte-for-byte what every reader already had"
);
{
let stage = super::repair_stage_partition(&partition, &fingerprints.before);
super::destroy_one(
&context,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&stage,
super::DestroyCheckpoint::Clear,
)
.await
.expect("clear the planted stage before the real run");
}
let scan = super::RepairScan {
readable: vec![survivor],
quarantined: vec![polyc_eventlog::QuarantinedItem {
position: 1,
error: "injected unreadable frame".to_owned(),
}],
before: inspected.before,
after,
standing: inspected.standing,
};
arm_append_failure(&partition, 0);
super::apply_repair_scan(
&context,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&partition,
scan,
None,
)
.await
.expect_err("the armed restore failure surfaces");
logs = super::OpenLogs::new();
checkpoints.clear();
count_checkpoints.evict_all();
mmr_logs.clear();
assert_eq!(
super::partition_event_count_one(&context, &run_dir, &mut logs, &partition)
.await
.unwrap(),
0,
"the interrupted repair leaves nothing readable — this is the state \
State used to serve as a legitimately empty conversation"
);
assert!(
!settled!(),
"a complete stage over a partition that does not match `after` is damage"
);
let (application, _) = super::apply_repair_conditionally(
&context,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&trust,
&partition,
fingerprints,
None,
)
.await
.expect("the resume finishes the repair from its stage");
assert_eq!(application, RepairApplication::Applied);
assert!(settled!(), "a finished repair is settled");
super::append_batch_one(
&context,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&partition,
&[Event::new("user_msg", b"after the repair".to_vec())],
)
.await
.expect("the repaired partition takes new commits");
assert!(
settled!(),
"a partition that grew past its repair is still settled"
);
});
let _ = std::fs::remove_dir_all(dir);
}
#[tokio::test]
#[allow(clippy::too_many_lines)]
async fn active_section_kind_corruption_is_a_hard_verify_failure_not_a_false_pass() {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-active-corrupt-kind-{}",
std::process::id()
));
let _ = std::fs::remove_dir_all(&dir);
let partition = "conv-corrupt-kind".to_owned();
{
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(
dir.clone(),
shutdown.clone(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
)
.expect("spawn");
host.append_batch(
partition.clone(),
vec![
Event::new("user_msg", b"hello".to_vec()),
Event::new("output_msg", b"hi there".to_vec()),
Event::new("tool_call", b"lookup".to_vec()),
Event::new("output_msg", b"tool result".to_vec()),
Event::new("user_msg", b"thanks".to_vec()),
Event::new("output_msg", b"you're welcome".to_vec()),
Event::new("turn_complete", Vec::new()),
],
)
.await
.expect("append 7-event turn");
let verified = host
.replay_verified(partition.clone())
.await
.expect("uncorrupted conversation verifies");
assert_eq!(verified.len(), 8, "7 content events + 1 signed-root marker");
drop(host);
}
let data_file = dir
.join(format!("{partition}_data"))
.join("0000000000000000");
let mut bytes = std::fs::read(&data_file).expect("read section 0");
let kind_at = bytes
.windows(b"user_msg".len())
.position(|w| w == b"user_msg")
.expect("the first event's kind bytes are present on disk");
bytes[kind_at] = 0xFF;
std::fs::write(&data_file, &bytes).expect("write corrupted section 0");
let shutdown2 = CancellationToken::new();
let host2 = EventLogHost::spawn(
dir.clone(),
shutdown2,
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
)
.expect("reopen");
let err = host2.replay_verified(partition.clone()).await.expect_err(
"a truncated replay of already-committed data must never report verified=true",
);
let AppendError::Verify(verify_err) = err else {
panic!("expected AppendError::Verify, got {err:?}");
};
match verify_err {
super::VerifyError::TruncatedReplay { expected, actual } => {
assert_eq!(expected, 8, "checkpoint recorded all 8 events as committed");
assert_eq!(
actual, 0,
"the journal's self-heal wiped the whole (only) active section"
);
}
other => panic!("expected TruncatedReplay, got {other:?}"),
}
let plan = host2
.inspect_repair_partition(partition.clone(), None)
.await
.expect("repair inspection completes");
let quarantined = plan.quarantined.clone();
assert_eq!(
quarantined.len(),
8,
"repair must flag the full gap between the checkpoint and what's visible, not report nothing"
);
let positions: Vec<u64> = quarantined.iter().map(|q| q.position).collect();
assert_eq!(positions, (0..8).collect::<Vec<_>>());
assert_eq!(
host2
.apply_repair_partition(partition.clone(), plan.before, plan.after, None)
.await
.expect("repair apply completes"),
RepairApplication::Applied
);
assert_eq!(
host2
.apply_repair_partition(partition.clone(), plan.before, plan.after, None)
.await
.expect("response-loss retry reconciles"),
RepairApplication::AlreadyApplied
);
let reverified = host2
.replay_verified(partition)
.await
.expect("post-repair verify judges the partition against its new, honest baseline");
assert!(
reverified.is_empty(),
"everything was unrecoverable in this reproduction; repair's job is honesty, not magic"
);
drop(host2);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn a_rewrite_that_crashed_before_lowering_the_count_is_healed() {
use commonware_runtime::{Runner as _, tokio as cw_tokio};
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-rewrite-count-window-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let _ = std::fs::remove_dir_all(&dir);
let partition = "conv-rewrite-count-window".to_owned();
let command_id = "excise-count-window".to_owned();
let run_dir = dir.clone();
let runner =
cw_tokio::Runner::new(cw_tokio::Config::default().with_storage_directory(dir.clone()));
runner.start(move |context| async move {
let mut logs = super::OpenLogs::new();
let mut checkpoints = HashMap::new();
let mut count_checkpoints = super::CountCheckpoints::new();
let mut mmr_logs = HashMap::new();
let signer = polyc_crypto::signing_role::JournalAttestationSigner::from_seed(29);
let trust = polyc_crypto::signing_role::RoleTrustSet::current(&signer);
let survivors = vec![Event::new("user_msg", b"kept".to_vec())];
let tree = polyc_eventlog::rebuild_from_events(&survivors).expect("tree");
let root = polyc_eventlog::extend_and_sign(&tree, &[], &signer).expect("sign");
let restored: Vec<Event> = survivors
.iter()
.cloned()
.chain(std::iter::once(root))
.collect();
{
let log = logs.get_or_open(&context, &partition).await.unwrap();
for event in &restored {
log.append(event).await.unwrap();
}
log.commit().await.unwrap();
}
count_checkpoints
.record(&context, &partition, 3)
.await
.unwrap();
let stage = super::rewrite_stage_partition(&partition, &command_id);
{
let mut staged = restored.clone();
staged.push(super::rewrite_stage_complete(restored.len() as u64, 1));
let log = logs.get_or_open(&context, &stage).await.unwrap();
for event in &staged {
log.append(event).await.unwrap();
}
log.commit().await.unwrap();
}
logs = super::OpenLogs::new();
count_checkpoints.evict_all();
let before = super::replay_verified_one(
&context,
&mut logs,
&mut count_checkpoints,
&trust,
&partition,
)
.await;
assert!(
matches!(before, Err(super::VerifyError::TruncatedReplay { .. })),
"the retained floor reports this correct partition as damaged: {before:?}"
);
let dropped = super::resume_rewrite(
&context,
&run_dir,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&partition,
&stage,
)
.await
.expect("the stage finishes")
.expect("a stage was waiting");
assert_eq!(dropped, 1);
let verified = super::replay_verified_one(
&context,
&mut logs,
&mut count_checkpoints,
&trust,
&partition,
)
.await
.expect("the healed partition verifies");
assert_eq!(
verified, restored,
"the content never changed, only the floor"
);
});
let _ = std::fs::remove_dir_all(&dir);
}
#[allow(clippy::too_many_lines)]
#[test]
fn a_staged_rewrite_is_finishable_by_a_later_unrelated_caller() {
use commonware_runtime::{Runner as _, tokio as cw_tokio};
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-rewrite-foreign-resume-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let _ = std::fs::remove_dir_all(&dir);
let partition = "conv-rewrite-foreign".to_owned();
let interrupted = "excise-interrupted".to_owned();
let run_dir = dir.clone();
let runner =
cw_tokio::Runner::new(cw_tokio::Config::default().with_storage_directory(dir.clone()));
runner.start(move |context| async move {
let mut logs = super::OpenLogs::new();
let mut checkpoints = HashMap::new();
let mut count_checkpoints = super::CountCheckpoints::new();
let mut mmr_logs = HashMap::new();
let signer = polyc_crypto::signing_role::JournalAttestationSigner::from_seed(23);
let content = vec![
Event::new("user_msg", b"alpha".to_vec()),
Event::new("secret", b"omega".to_vec()),
];
let tree = polyc_eventlog::rebuild_from_events(&content).expect("tree");
let root = polyc_eventlog::extend_and_sign(&tree, &[], &signer).expect("sign");
{
let log = logs.get_or_open(&context, &partition).await.unwrap();
for event in content.iter().chain(std::iter::once(&root)) {
log.append(event).await.unwrap();
}
log.commit().await.unwrap();
}
count_checkpoints
.record(&context, &partition, 3)
.await
.unwrap();
arm_append_failure(&partition, 0);
super::rewrite_one(
&context,
&run_dir,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&interrupted,
&partition,
&|event: &Event| {
if event.kind == "secret" {
RewriteDecision::Drop
} else {
RewriteDecision::Keep
}
},
)
.await
.expect_err("the armed restore failure surfaces");
logs = super::OpenLogs::new();
checkpoints.clear();
count_checkpoints.evict_all();
mmr_logs.clear();
let stage = super::rewrite_stage_partition(&partition, &interrupted);
let dropped = super::resume_rewrite(
&context,
&run_dir,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&partition,
&stage,
)
.await
.expect("the stage finishes")
.expect("a stage was waiting");
assert_eq!(dropped, 1);
let restored = logs
.get_or_open(&context, &partition)
.await
.unwrap()
.replay()
.await
.unwrap();
assert_eq!(restored.len(), 2, "one survivor plus one fresh root");
assert_eq!(restored[0], content[0]);
let second = super::resume_rewrite(
&context,
&run_dir,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&partition,
&stage,
)
.await
.expect("a second resume is a no-op");
assert!(
second.is_none(),
"a finished excision leaves no stage, so recovery costs one call"
);
});
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
#[allow(clippy::too_many_lines)]
fn an_interrupted_rewrite_fails_loud_and_resumes_from_its_stage() {
use commonware_runtime::{Runner as _, tokio as cw_tokio};
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-rewrite-resume-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let _ = std::fs::remove_dir_all(&dir);
let partition = "conv-rewrite-resume".to_owned();
let command_id = "excise-resume-1".to_owned();
let run_dir = dir.clone();
let runner =
cw_tokio::Runner::new(cw_tokio::Config::default().with_storage_directory(dir.clone()));
runner.start(move |context| async move {
let mut logs = super::OpenLogs::new();
let mut checkpoints = HashMap::new();
let mut count_checkpoints = super::CountCheckpoints::new();
let mut mmr_logs = HashMap::new();
let signer = polyc_crypto::signing_role::JournalAttestationSigner::from_seed(21);
let trust = polyc_crypto::signing_role::RoleTrustSet::current(&signer);
let content = vec![
Event::new("user_msg", b"keep-one".to_vec()),
Event::new("secret", b"drop-me".to_vec()),
Event::new("user_msg", b"keep-two".to_vec()),
];
let tree = polyc_eventlog::rebuild_from_events(&content).expect("tree");
let root = polyc_eventlog::extend_and_sign(&tree, &[], &signer).expect("sign");
let seeded: Vec<Event> = content
.iter()
.cloned()
.chain(std::iter::once(root))
.collect();
{
let log = logs.get_or_open(&context, &partition).await.unwrap();
for event in &seeded {
log.append(event).await.unwrap();
}
log.commit().await.unwrap();
}
count_checkpoints
.record(&context, &partition, seeded.len() as u64)
.await
.unwrap();
let decide = |event: &Event| {
if event.kind == "secret" {
RewriteDecision::Drop
} else {
RewriteDecision::Keep
}
};
arm_append_failure(&partition, 0);
super::rewrite_one(
&context,
&run_dir,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&command_id,
&partition,
&decide,
)
.await
.expect_err("the armed restore failure surfaces");
logs = super::OpenLogs::new();
checkpoints.clear();
count_checkpoints.evict_all();
mmr_logs.clear();
let verdict = super::replay_verified_one(
&context,
&mut logs,
&mut count_checkpoints,
&trust,
&partition,
)
.await;
match verdict {
Err(super::VerifyError::TruncatedReplay { expected, actual }) => {
assert_eq!(expected, 4, "the pre-rewrite count is still the floor");
assert_eq!(actual, 0, "the destroy really did empty the partition");
}
other => panic!("an interrupted rewrite must report damage, got {other:?}"),
}
let dropped = super::rewrite_one(
&context,
&run_dir,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&command_id,
&partition,
&decide,
)
.await
.expect("the retry resumes from the durable stage");
assert_eq!(
dropped, 1,
"the resume reports what the interrupted run dropped"
);
let restored = logs
.get_or_open(&context, &partition)
.await
.unwrap()
.replay()
.await
.unwrap();
assert_eq!(
restored.len(),
3,
"two surviving records plus one fresh root"
);
assert_eq!(restored[0], content[0]);
assert_eq!(restored[1], content[2]);
assert_eq!(restored[2].kind, polyc_eventlog::MMR_SIGNED_ROOT_KIND);
let verified = super::replay_verified_one(
&context,
&mut logs,
&mut count_checkpoints,
&trust,
&partition,
)
.await
.expect("a resumed rewrite leaves a verifying partition");
assert_eq!(verified.len(), 3);
assert!(
super::list_partitions_in(&dir)
.unwrap()
.iter()
.all(|name| !name.starts_with(super::REWRITE_STAGE_PREFIX)),
"rewrite staging is never a caller-visible partition"
);
});
}
#[test]
#[allow(clippy::too_many_lines)]
fn a_rewrite_whose_terminal_count_record_fails_keeps_its_stage() {
use commonware_runtime::{Runner as _, tokio as cw_tokio};
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-rewrite-record-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let _ = std::fs::remove_dir_all(&dir);
let partition = "conv-rewrite-record".to_owned();
let command_id = "excise-record-1".to_owned();
let run_dir = dir.clone();
let runner =
cw_tokio::Runner::new(cw_tokio::Config::default().with_storage_directory(dir.clone()));
runner.start(move |context| async move {
let mut logs = super::OpenLogs::new();
let mut checkpoints = HashMap::new();
let mut count_checkpoints = super::CountCheckpoints::new();
let mut mmr_logs = HashMap::new();
let signer = polyc_crypto::signing_role::JournalAttestationSigner::from_seed(31);
let trust = polyc_crypto::signing_role::RoleTrustSet::current(&signer);
let content = vec![
Event::new("user_msg", b"keep-one".to_vec()),
Event::new("secret", b"drop-me".to_vec()),
Event::new("user_msg", b"keep-two".to_vec()),
];
let tree = polyc_eventlog::rebuild_from_events(&content).expect("tree");
let root = polyc_eventlog::extend_and_sign(&tree, &[], &signer).expect("sign");
let seeded: Vec<Event> = content
.iter()
.cloned()
.chain(std::iter::once(root))
.collect();
{
let log = logs.get_or_open(&context, &partition).await.unwrap();
for event in &seeded {
log.append(event).await.unwrap();
}
log.commit().await.unwrap();
}
count_checkpoints
.record(&context, &partition, seeded.len() as u64)
.await
.unwrap();
let decide = |event: &Event| {
if event.kind == "secret" {
RewriteDecision::Drop
} else {
RewriteDecision::Keep
}
};
arm_count_record_failure(&partition);
super::rewrite_one(
&context,
&run_dir,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&command_id,
&partition,
&decide,
)
.await
.expect_err("the terminal count record's failure must surface");
logs = super::OpenLogs::new();
checkpoints.clear();
count_checkpoints.evict_all();
mmr_logs.clear();
match super::replay_verified_one(
&context,
&mut logs,
&mut count_checkpoints,
&trust,
&partition,
)
.await
{
Err(super::VerifyError::TruncatedReplay { expected, actual }) => {
assert_eq!(expected, 4, "the pre-rewrite count is still the floor");
assert_eq!(actual, 3, "two survivors plus one fresh root did land");
}
other => panic!("a rewrite short of its record must report damage, got {other:?}"),
}
assert!(
stage_dirs_in(&run_dir) > 0,
"the stage must outlive a failed terminal record"
);
let dropped = super::rewrite_one(
&context,
&run_dir,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&command_id,
&partition,
&decide,
)
.await
.expect("the retry finishes the rewrite from its stage");
assert_eq!(
dropped, 1,
"the resume reports what the interrupted run dropped"
);
let verified = super::replay_verified_one(
&context,
&mut logs,
&mut count_checkpoints,
&trust,
&partition,
)
.await
.expect("the resumed rewrite leaves a partition that reads clean");
assert_eq!(verified.len(), 3, "two survivors plus one fresh root");
assert_eq!(verified[0], content[0]);
assert_eq!(verified[1], content[2]);
assert_eq!(
stage_dirs_in(&run_dir),
0,
"the finished rewrite removes its private stage"
);
});
let _ = std::fs::remove_dir_all(dir);
}
fn stage_dirs_in(storage_dir: &std::path::Path) -> usize {
std::fs::read_dir(storage_dir)
.expect("the storage directory exists")
.filter_map(Result::ok)
.filter(|entry| {
let name = entry.file_name().to_string_lossy().into_owned();
name.starts_with(super::REWRITE_STAGE_PREFIX)
&& !name.contains("__eventcount_checkpoint")
})
.count()
}
#[test]
#[allow(clippy::too_many_lines)]
fn a_resumed_repair_re_roots_and_leaves_an_extendable_tree() {
use commonware_runtime::{Runner as _, tokio as cw_tokio};
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-repair-resume-reroot-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let _ = std::fs::remove_dir_all(&dir);
let partition = "conv-repair-resume-reroot".to_owned();
let runner =
cw_tokio::Runner::new(cw_tokio::Config::default().with_storage_directory(dir.clone()));
runner.start(move |context| async move {
let mut logs = super::OpenLogs::new();
let mut checkpoints = HashMap::new();
let mut count_checkpoints = super::CountCheckpoints::new();
let mut mmr_logs = HashMap::new();
let signer = polyc_crypto::signing_role::JournalAttestationSigner::from_seed(11);
let trust = polyc_crypto::signing_role::RoleTrustSet::current(&signer);
let content = vec![
Event::new("user_msg", b"alpha".to_vec()),
Event::new("output_msg", b"beta".to_vec()),
Event::new("tool_call", b"gamma".to_vec()),
];
let tree = polyc_eventlog::rebuild_from_events(&content).expect("tree");
let stale_root = polyc_eventlog::extend_and_sign(&tree, &[], &signer).expect("sign");
let seeded: Vec<Event> = content
.iter()
.cloned()
.chain(std::iter::once(stale_root.clone()))
.collect();
{
let log = logs.get_or_open(&context, &partition).await.unwrap();
for event in &seeded {
log.append(event).await.unwrap();
}
log.commit().await.unwrap();
}
count_checkpoints
.record(&context, &partition, seeded.len() as u64)
.await
.unwrap();
let readable: Vec<(u64, Event)> = seeded
.iter()
.enumerate()
.filter(|(position, _)| *position != 1)
.map(|(position, event)| (position as u64, event.clone()))
.collect();
let quarantined = vec![polyc_eventlog::QuarantinedItem {
position: 1,
error: "injected unreadable frame".to_owned(),
}];
let raw: Vec<Event> = readable.iter().map(|(_, event)| event.clone()).collect();
let (expected_content, _) =
super::repaired_content(&raw, &signer, None).expect("re-root");
let indexed: Vec<(u64, Event)> = expected_content
.iter()
.enumerate()
.map(|(position, event)| (position as u64, event.clone()))
.collect();
let before = super::repair_fingerprint(&readable, &quarantined, seeded.len() as u64);
let after = super::repair_fingerprint(&indexed, &[], expected_content.len() as u64);
let scan = super::RepairScan {
readable,
quarantined,
before,
after,
standing: polyc_eventlog::RootStanding::Stale,
};
arm_append_failure(&partition, 0);
super::apply_repair_scan(
&context,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&partition,
scan,
None,
)
.await
.expect_err("the armed restore failure surfaces");
logs = super::OpenLogs::new();
checkpoints.clear();
count_checkpoints.evict_all();
mmr_logs.clear();
let (application, _) = super::apply_repair_conditionally(
&context,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&trust,
&partition,
super::RepairFingerprints { before, after },
None,
)
.await
.expect("the restart resumes from the durable stage");
assert_eq!(application, RepairApplication::Applied);
let restored = logs
.get_or_open(&context, &partition)
.await
.unwrap()
.replay()
.await
.unwrap();
assert_eq!(
restored, expected_content,
"a resumed repair restores the re-rooted content"
);
let roots: Vec<&Event> = restored
.iter()
.filter(|event| event.kind == polyc_eventlog::MMR_SIGNED_ROOT_KIND)
.collect();
assert_eq!(roots.len(), 1, "exactly one root survives the resume");
assert_ne!(*roots[0], stale_root, "and it is not the stale one");
super::replay_verified_one(
&context,
&mut logs,
&mut count_checkpoints,
&trust,
&partition,
)
.await
.expect("a resumed repair leaves a verifying partition");
super::append_batch_one(
&context,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&partition,
&[Event::new("user_msg", b"delta".to_vec())],
)
.await
.expect("the repaired partition accepts another batch");
let verified = super::replay_verified_one(
&context,
&mut logs,
&mut count_checkpoints,
&trust,
&partition,
)
.await
.expect("the partition still verifies after extending it");
assert!(
verified.len() > restored.len(),
"the extending batch landed on top of the repaired content"
);
});
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
#[allow(clippy::too_many_lines)]
fn a_repair_that_drops_a_position_re_roots_what_survives() {
use commonware_runtime::{Runner as _, tokio as cw_tokio};
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-repair-reroot-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let _ = std::fs::remove_dir_all(&dir);
let partition = "conv-repair-reroot".to_owned();
let runner =
cw_tokio::Runner::new(cw_tokio::Config::default().with_storage_directory(dir.clone()));
runner.start(move |context| async move {
let mut logs = super::OpenLogs::new();
let mut checkpoints = HashMap::new();
let mut count_checkpoints = super::CountCheckpoints::new();
let mut mmr_logs = HashMap::new();
let signer = polyc_crypto::signing_role::JournalAttestationSigner::from_seed(9);
let trust = polyc_crypto::signing_role::RoleTrustSet::current(&signer);
let content = vec![
Event::new("user_msg", b"one".to_vec()),
Event::new("output_msg", b"two".to_vec()),
Event::new("tool_call", b"three".to_vec()),
Event::new("output_msg", b"four".to_vec()),
];
let tree = polyc_eventlog::rebuild_from_events(&content).expect("tree");
let original_root = polyc_eventlog::extend_and_sign(&tree, &[], &signer).expect("sign");
let seeded: Vec<Event> = content
.iter()
.cloned()
.chain(std::iter::once(original_root.clone()))
.collect();
{
let log = logs.get_or_open(&context, &partition).await.unwrap();
for event in &seeded {
log.append(event).await.unwrap();
}
log.commit().await.unwrap();
}
count_checkpoints
.record(&context, &partition, seeded.len() as u64)
.await
.unwrap();
let readable: Vec<(u64, Event)> = seeded
.iter()
.enumerate()
.filter(|(position, _)| *position != 2)
.map(|(position, event)| (position as u64, event.clone()))
.collect();
assert!(
readable
.iter()
.any(|(_, event)| event.kind == polyc_eventlog::MMR_SIGNED_ROOT_KIND),
"the fixture must carry a stale root into the repair"
);
let quarantined = vec![polyc_eventlog::QuarantinedItem {
position: 2,
error: "injected unreadable frame".to_owned(),
}];
let raw: Vec<Event> = readable.iter().map(|(_, event)| event.clone()).collect();
let (expected_content, _) =
super::repaired_content(&raw, &signer, None).expect("re-root");
let indexed: Vec<(u64, Event)> = expected_content
.iter()
.enumerate()
.map(|(position, event)| (position as u64, event.clone()))
.collect();
let scan = super::RepairScan {
before: super::repair_fingerprint(&readable, &quarantined, seeded.len() as u64),
after: super::repair_fingerprint(&indexed, &[], expected_content.len() as u64),
readable,
quarantined,
standing: polyc_eventlog::RootStanding::Stale,
};
super::apply_repair_scan(
&context,
&mut logs,
&mut checkpoints,
&mut count_checkpoints,
&mut mmr_logs,
&signer,
&partition,
scan,
None,
)
.await
.expect("the repair applies");
let restored = logs
.get_or_open(&context, &partition)
.await
.unwrap()
.replay()
.await
.unwrap();
let roots: Vec<&Event> = restored
.iter()
.filter(|event| event.kind == polyc_eventlog::MMR_SIGNED_ROOT_KIND)
.collect();
assert_eq!(roots.len(), 1, "exactly one root survives a repair");
assert_ne!(
*roots[0], original_root,
"the surviving root is freshly signed, not the stale one"
);
assert_eq!(
restored.len(),
4,
"three surviving content events plus one fresh root"
);
let verified = super::replay_verified_one(
&context,
&mut logs,
&mut count_checkpoints,
&trust,
&partition,
)
.await
.expect("a repaired partition verifies against its own fresh root");
assert_eq!(verified.len(), 4);
});
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn active_section_payload_content_tamper_is_a_hard_verify_failure_not_a_false_pass() {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-active-corrupt-payload-{}",
std::process::id()
));
let _ = std::fs::remove_dir_all(&dir);
let partition = "conv-corrupt-payload".to_owned();
{
let shutdown = CancellationToken::new();
let host = EventLogHost::spawn(
dir.clone(),
shutdown.clone(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
)
.expect("spawn");
host.append_batch(
partition.clone(),
vec![
Event::new("user_msg", b"what is 2+2?".to_vec()),
Event::new("output_msg", b"the reply is four".to_vec()),
],
)
.await
.expect("append 2-event turn");
let verified = host
.replay_verified(partition.clone())
.await
.expect("uncorrupted conversation verifies");
assert_eq!(verified.len(), 3, "2 content events + 1 signed-root marker");
drop(host);
}
let data_file = dir
.join(format!("{partition}_data"))
.join("0000000000000000");
let mut bytes = std::fs::read(&data_file).expect("read section 0");
let payload_at = bytes
.windows(b"the reply is four".len())
.position(|w| w == b"the reply is four")
.expect("the second event's payload bytes are present on disk");
bytes[payload_at] ^= 0xFF;
std::fs::write(&data_file, &bytes).expect("write corrupted section 0");
let shutdown2 = CancellationToken::new();
let host2 = EventLogHost::spawn(
dir.clone(),
shutdown2,
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
)
.expect("reopen");
let err = host2
.replay_verified(partition.clone())
.await
.expect_err("a content-only tamper must still be caught, not silently truncated");
let AppendError::Verify(verify_err) = err else {
panic!("expected AppendError::Verify, got {err:?}");
};
match verify_err {
super::VerifyError::TruncatedReplay { expected, actual } => {
assert_eq!(expected, 3, "checkpoint recorded all 3 events as committed");
assert_eq!(
actual, 0,
"the storage layer's active-section revalidation invalidates the WHOLE \
section on any internal inconsistency, not just a truncated tail — even \
the untouched first event is lost"
);
}
other => panic!("expected TruncatedReplay, got {other:?}"),
}
let quarantined = host2
.repair_partition(partition.clone())
.await
.expect("repair completes");
assert_eq!(
quarantined.len(),
3,
"repair must flag the full gap between the checkpoint and what's visible, not report nothing"
);
let positions: Vec<u64> = quarantined.iter().map(|q| q.position).collect();
assert_eq!(positions, vec![0, 1, 2]);
let reverified = host2
.replay_verified(partition)
.await
.expect("post-repair verify judges the partition against its new, honest baseline");
assert!(
reverified.is_empty(),
"everything was unrecoverable in this reproduction; repair's job is honesty, not magic"
);
drop(host2);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test(flavor = "multi_thread")]
async fn shard_slow_replay_on_one_partition_does_not_block_appends_on_another() {
const SLOW_PARTITION_EVENTS: usize = 20_000;
let dir =
std::env::temp_dir().join(format!("polychrome-eventlog-shard-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let shutdown = CancellationToken::new();
let host = std::sync::Arc::new(
EventLogHost::spawn(
dir.clone(),
shutdown.clone(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
)
.expect("spawn"),
);
assert!(
host.num_shards() > 1,
"sharding requires more than one shard"
);
let mut partition_a = None;
let mut partition_b = None;
for i in 0..64u32 {
let name = format!("conv-shard-{i}");
let idx = host.shard_index_for(&name);
match (&partition_a, &partition_b) {
(None, _) => partition_a = Some((name, idx)),
(Some((_, idx_a)), None) if idx != *idx_a => partition_b = Some((name, idx)),
_ => {}
}
}
let (partition_a, _) = partition_a.expect("found a first partition");
let (partition_b, _) = partition_b.expect("found a partition on a different shard");
let mut batch = Vec::with_capacity(SLOW_PARTITION_EVENTS);
for i in 0..SLOW_PARTITION_EVENTS {
batch.push(Event::new(format!("k{i}"), vec![b'x'; 512]));
}
host.append_batch(partition_a.clone(), batch)
.await
.expect("seed the slow partition");
let replay_host = host.clone();
let replay_partition = partition_a.clone();
let replay_handle = tokio::spawn(async move { replay_host.replay(replay_partition).await });
let start = std::time::Instant::now();
for i in 0..20u32 {
host.append_batch(
partition_b.clone(),
vec![Event::new(format!("b{i}"), b"fast".to_vec())],
)
.await
.expect("append to the other shard");
}
let elapsed = start.elapsed();
replay_handle
.await
.expect("replay task")
.expect("slow replay itself succeeds");
let content_events = host
.replay(partition_b)
.await
.expect("replay b")
.into_iter()
.filter(|e| e.kind != polyc_eventlog::MMR_SIGNED_ROOT_KIND)
.count();
assert_eq!(content_events, 20, "all 20 fast appends landed");
assert!(
elapsed < std::time::Duration::from_secs(5),
"appends to a different shard took {elapsed:?} — too slow to plausibly be \
running concurrently with the other shard's big replay"
);
drop(host);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test(flavor = "multi_thread")]
async fn a_paused_replay_on_one_partition_does_not_block_appends_on_another_partition_on_the_same_shard()
{
const FAST_APPENDS: u32 = 4;
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-same-shard-{}",
std::process::id()
));
let _ = std::fs::remove_dir_all(&dir);
let shutdown = CancellationToken::new();
let host = std::sync::Arc::new(
EventLogHost::spawn(
dir.clone(),
shutdown.clone(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
)
.expect("spawn"),
);
let mut partition_a = None;
let mut partition_b = None;
for i in 0..64u32 {
let name = format!("conv-same-shard-{i}");
let idx = host.shard_index_for(&name);
match (&partition_a, &partition_b) {
(None, _) => partition_a = Some((name, idx)),
(Some((_, idx_a)), None) if idx == *idx_a => partition_b = Some((name, idx)),
_ => {}
}
}
let (partition_a, _) = partition_a.expect("found a first partition");
let (partition_b, _) = partition_b.expect("found a partition on the SAME shard");
host.append_batch(
partition_a.clone(),
vec![Event::new("seed", b"one".to_vec())],
)
.await
.expect("seed the paused partition");
let pause = arm_replay_pause(&partition_a);
let replay_host = host.clone();
let replay_partition = partition_a.clone();
let replay_handle = tokio::spawn(async move { replay_host.replay(replay_partition).await });
tokio::time::timeout(std::time::Duration::from_secs(30), pause.entered.notified())
.await
.expect("the replay reached its deterministic pause");
let append_result = tokio::time::timeout(std::time::Duration::from_secs(30), async {
for i in 0..FAST_APPENDS {
host.append_batch(
partition_b.clone(),
vec![Event::new(format!("b{i}"), b"fast".to_vec())],
)
.await
.expect("append to the other partition on the same shard");
}
})
.await;
pause.release.notify_one();
replay_handle
.await
.expect("replay task")
.expect("slow replay itself succeeds");
let content_events = host
.replay(partition_b)
.await
.expect("replay b")
.into_iter()
.filter(|e| e.kind != polyc_eventlog::MMR_SIGNED_ROOT_KIND)
.count();
assert_eq!(
content_events, FAST_APPENDS as usize,
"all fast appends landed"
);
append_result.expect(
"appends to another partition on the SAME shard must complete while the first \
partition's replay is deliberately paused (#1711)",
);
drop(host);
let _ = std::fs::remove_dir_all(&dir);
}
}
#[cfg(test)]
mod presence_tests {
#![allow(
clippy::pedantic,
clippy::nursery,
missing_docs,
reason = "a test module: panics are the failure mode"
)]
use super::*;
use tokio_util::sync::CancellationToken;
fn scratch(name: &str) -> std::path::PathBuf {
let dir = std::env::temp_dir().join(format!(
"polyc-presence-{name}-{}-{}",
std::process::id(),
uuid::Uuid::now_v7().as_simple()
));
let _ = std::fs::remove_dir_all(&dir);
dir
}
fn spawn_host(dir: &std::path::Path, shutdown: &CancellationToken) -> EventLogHost {
EventLogHost::spawn(
dir.to_path_buf(),
shutdown.clone(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
)
.expect("spawn the host")
}
#[tokio::test]
async fn a_never_written_partition_is_absent_and_asking_creates_nothing() {
let dir = scratch("never-written");
std::fs::create_dir_all(&dir).unwrap();
let shutdown = CancellationToken::new();
let host = spawn_host(&dir, &shutdown);
assert_eq!(
host.partition_presence("conv-typo".to_owned())
.await
.unwrap(),
PartitionPresence::Absent
);
assert!(
std::fs::read_dir(&dir).unwrap().next().is_none(),
"asking about a typo must leave the volume empty"
);
shutdown.cancel();
let _ = std::fs::remove_dir_all(dir);
}
#[tokio::test]
async fn a_destroyed_partition_is_absent_after_its_checkpoint_is_cleared() {
let dir = scratch("destroyed");
let shutdown = CancellationToken::new();
let host = spawn_host(&dir, &shutdown);
host.append_batch("conv-gone".to_owned(), vec![Event::new("k", b"v".to_vec())])
.await
.unwrap();
assert_eq!(
host.partition_presence("conv-gone".to_owned())
.await
.unwrap(),
PartitionPresence::Present
);
host.destroy_partition("conv-gone".to_owned())
.await
.unwrap();
assert!(
dir.join("conv-gone__eventcount_checkpoint").exists(),
"a destroy clears the checkpoint in place rather than removing it"
);
assert_eq!(
host.partition_presence("conv-gone".to_owned())
.await
.unwrap(),
PartitionPresence::Absent
);
shutdown.cancel();
let _ = std::fs::remove_dir_all(dir);
}
#[tokio::test]
async fn a_lost_data_directory_with_a_surviving_index_is_damaged() {
let dir = scratch("lost-data");
let shutdown = CancellationToken::new();
let host = spawn_host(&dir, &shutdown);
host.append_batch("conv-lost".to_owned(), vec![Event::new("k", b"v".to_vec())])
.await
.unwrap();
shutdown.cancel();
drop(host);
std::fs::remove_dir_all(dir.join("conv-lost_data")).unwrap();
assert!(dir.join("conv-lost_offsets-blobs").exists());
let shutdown = CancellationToken::new();
let host = spawn_host(&dir, &shutdown);
match host
.partition_presence("conv-lost".to_owned())
.await
.unwrap()
{
PartitionPresence::Damaged { remnant } => {
assert_eq!(remnant, "conv-lost_offsets-blobs");
}
other => panic!("a lost data directory is damage, got: {other:?}"),
}
shutdown.cancel();
let _ = std::fs::remove_dir_all(dir);
}
#[tokio::test]
async fn a_surviving_nonzero_count_checkpoint_without_journal_data_is_damaged() {
let dir = scratch("lost-all-but-count");
let shutdown = CancellationToken::new();
let host = spawn_host(&dir, &shutdown);
host.append_batch(
"conv-count".to_owned(),
vec![Event::new("k", b"v".to_vec())],
)
.await
.unwrap();
shutdown.cancel();
drop(host);
std::fs::remove_dir_all(dir.join("conv-count_data")).unwrap();
std::fs::remove_dir_all(dir.join("conv-count_offsets-blobs")).unwrap();
std::fs::remove_dir_all(dir.join("conv-count_offsets-metadata")).unwrap();
let shutdown = CancellationToken::new();
let host = spawn_host(&dir, &shutdown);
match host
.partition_presence("conv-count".to_owned())
.await
.unwrap()
{
PartitionPresence::Damaged { remnant } => {
assert_eq!(remnant, "conv-count__eventcount_checkpoint");
}
other => panic!("a surviving nonzero count is damage, got: {other:?}"),
}
shutdown.cancel();
let _ = std::fs::remove_dir_all(dir);
}
#[tokio::test]
async fn an_error_on_the_first_event_is_unknown_too() {
let dir = scratch("first-event-unknown");
let shutdown = CancellationToken::new();
let host = spawn_host(&dir, &shutdown);
crate::tests::arm_append_failure("conv-first-unknown", 0);
let error = host
.append_batch(
"conv-first-unknown".to_owned(),
vec![
Event::new("output_msg", b"only".to_vec()),
Event::new("output_msg", b"second".to_vec()),
],
)
.await
.expect_err("the first event was armed to fail");
assert!(
matches!(error, AppendError::AppendOutcomeUnknown(_)),
"an errored append proves nothing about what it wrote: {error:?}"
);
shutdown.cancel();
let _ = std::fs::remove_dir_all(dir);
}
#[test]
fn will_rewrite_follows_the_root_standing_not_the_quarantine_list() {
let scan = |quarantined: Vec<polyc_eventlog::QuarantinedItem>,
standing: polyc_eventlog::RootStanding,
readable: Vec<(u64, Event)>| RepairScan {
readable,
quarantined,
before: [1; 32],
after: [2; 32],
standing,
};
let one = || vec![(0_u64, Event::new("user_msg", b"alpha".to_vec()))];
let quarantined = || {
vec![polyc_eventlog::QuarantinedItem {
position: 1,
error: "unreadable".to_owned(),
}]
};
assert!(
scan(Vec::new(), polyc_eventlog::RootStanding::Stale, one()).will_rewrite(),
"a stale root is repair work even with nothing to quarantine"
);
assert!(
scan(quarantined(), polyc_eventlog::RootStanding::Holds, one()).will_rewrite(),
"quarantining a record rewrites the partition"
);
assert!(
!scan(Vec::new(), polyc_eventlog::RootStanding::Holds, Vec::new()).will_rewrite(),
"an empty partition holds, and has no marker a replacement could rewrite"
);
}
#[tokio::test]
async fn probing_one_damaged_partition_twice_answers_the_same() {
let dir = scratch("probe-twice");
let shutdown = CancellationToken::new();
let host = spawn_host(&dir, &shutdown);
host.append_batch(
"conv-twice".to_owned(),
vec![Event::new("k", b"v".to_vec())],
)
.await
.unwrap();
shutdown.cancel();
drop(host);
std::fs::remove_dir_all(dir.join("conv-twice_data")).unwrap();
std::fs::remove_dir_all(dir.join("conv-twice_offsets-blobs")).unwrap();
std::fs::remove_dir_all(dir.join("conv-twice_offsets-metadata")).unwrap();
let shutdown = CancellationToken::new();
let host = spawn_host(&dir, &shutdown);
let first = host
.partition_presence("conv-twice".to_owned())
.await
.unwrap();
let second = host
.partition_presence("conv-twice".to_owned())
.await
.unwrap();
assert_eq!(first, second, "a probe is a read and does not change state");
assert!(
matches!(first, PartitionPresence::Damaged { .. }),
"got: {first:?}"
);
shutdown.cancel();
let _ = std::fs::remove_dir_all(dir);
}
#[tokio::test]
async fn discovery_refuses_a_partition_shaped_remnant_rather_than_omitting_it() {
let dir = scratch("listing-remnant");
let shutdown = CancellationToken::new();
let host = spawn_host(&dir, &shutdown);
host.append_batch("conv-keep".to_owned(), vec![Event::new("k", b"v".to_vec())])
.await
.unwrap();
host.append_batch("conv-half".to_owned(), vec![Event::new("k", b"v".to_vec())])
.await
.unwrap();
shutdown.cancel();
drop(host);
std::fs::remove_dir_all(dir.join("conv-half_data")).unwrap();
match crate::list_partitions_bounded_in(&dir, usize::MAX, usize::MAX, None) {
Err(ListPartitionsError::Corrupt { entry }) => {
assert!(entry.starts_with("conv-half"), "got: {entry}");
}
other => panic!("discovery must fail whole on a remnant, got: {other:?}"),
}
let _ = std::fs::remove_dir_all(dir);
}
#[tokio::test]
async fn an_excluded_partition_does_not_look_like_a_remnant() {
let dir = scratch("listing-excluded");
let shutdown = CancellationToken::new();
let host = spawn_host(&dir, &shutdown);
host.append_batch("conv-keep".to_owned(), vec![Event::new("k", b"v".to_vec())])
.await
.unwrap();
host.append_batch(
"conv-hidden".to_owned(),
vec![Event::new("k", b"v".to_vec())],
)
.await
.unwrap();
shutdown.cancel();
drop(host);
assert!(dir.join("conv-hidden_offsets-blobs").exists());
let listed =
crate::list_partitions_bounded_in(&dir, usize::MAX, usize::MAX, Some("conv-hidden"))
.expect("an excluded partition is not a remnant");
assert_eq!(listed, vec!["conv-keep".to_owned()]);
let _ = std::fs::remove_dir_all(dir);
}
#[tokio::test]
async fn discovery_accepts_the_cleared_checkpoint_a_destroy_leaves() {
let dir = scratch("listing-destroyed");
let shutdown = CancellationToken::new();
let host = spawn_host(&dir, &shutdown);
host.append_batch("conv-keep".to_owned(), vec![Event::new("k", b"v".to_vec())])
.await
.unwrap();
host.append_batch("conv-gone".to_owned(), vec![Event::new("k", b"v".to_vec())])
.await
.unwrap();
host.destroy_partition("conv-gone".to_owned())
.await
.unwrap();
shutdown.cancel();
drop(host);
let listed =
crate::list_partitions_bounded_in(&dir, usize::MAX, usize::MAX, None).expect("lists");
assert_eq!(listed, vec!["conv-keep".to_owned()]);
let _ = std::fs::remove_dir_all(dir);
}
}