use std::collections::HashMap;
use std::panic::AssertUnwindSafe;
use std::sync::{Arc, Mutex, RwLock};
use polyc_eventlog::Event;
type ObserverList = Arc<RwLock<Vec<Arc<dyn EventLogObserver>>>>;
type EpochMap = Arc<Mutex<HashMap<String, u64>>>;
#[derive(Clone)]
pub(crate) struct WriteObservers {
observers: ObserverList,
epochs: EpochMap,
}
impl WriteObservers {
pub(crate) fn new() -> Self {
Self {
observers: Arc::new(RwLock::new(Vec::new())),
epochs: Arc::new(Mutex::new(HashMap::new())),
}
}
pub(crate) fn register(&self, observer: Arc<dyn EventLogObserver>) {
self.observers.write().expect("poison").push(observer);
}
pub(crate) fn epoch(&self, partition: &str) -> u64 {
self.epochs
.lock()
.expect("poison")
.get(partition)
.copied()
.unwrap_or(0)
}
pub(crate) fn notify_append(&self, partition: &str, events: &[Event], positions: &[u64]) {
if events.is_empty() {
return;
}
let observers = {
let guard = self.observers.read().expect("poison");
if guard.is_empty() {
return;
}
guard.clone()
};
let notification = AppendNotification {
partition,
events,
positions,
};
for observer in &observers {
if let Err(payload) =
std::panic::catch_unwind(AssertUnwindSafe(|| observer.on_append(¬ification)))
{
log_observer_panic("on_append", partition, &*payload);
}
}
}
pub(crate) fn notify_mutation_bumping_epoch(
&self,
partition: &str,
kind: &MutationKind,
) -> u64 {
let mut guard = self.epochs.lock().expect("poison");
let epoch = guard.entry(partition.to_owned()).or_insert(0);
*epoch += 1;
let epoch = *epoch;
drop(guard);
let observers = {
let guard = self.observers.read().expect("poison");
if guard.is_empty() {
return epoch;
}
guard.clone()
};
let notification = MutationNotification {
partition,
kind,
epoch,
};
for observer in &observers {
if let Err(payload) = std::panic::catch_unwind(AssertUnwindSafe(|| {
observer.on_mutation(¬ification);
})) {
log_observer_panic("on_mutation", partition, &*payload);
}
}
epoch
}
}
pub trait EventLogObserver: Send + Sync {
fn on_append(&self, _notification: &AppendNotification<'_>) {}
fn on_mutation(&self, _notification: &MutationNotification<'_>) {}
}
#[derive(Debug, Clone, Copy)]
pub struct AppendNotification<'a> {
pub partition: &'a str,
pub events: &'a [Event],
pub positions: &'a [u64],
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum MutationKind {
Destroyed,
Rewritten {
dropped: usize,
},
MigrationSource {
to: String,
},
MigrationDestination {
from: String,
appended: usize,
},
Repaired {
quarantined: usize,
},
}
#[derive(Debug, Clone, Copy)]
pub struct MutationNotification<'a> {
pub partition: &'a str,
pub kind: &'a MutationKind,
pub epoch: u64,
}
fn log_observer_panic(hook: &str, partition: &str, payload: &(dyn std::any::Any + Send)) {
let message = payload
.downcast_ref::<&str>()
.map(|s| (*s).to_owned())
.or_else(|| payload.downcast_ref::<String>().cloned())
.unwrap_or_else(|| "non-string panic payload".to_owned());
tracing::error!(
hook,
partition = %partition,
message,
"event-log observer panicked; the write it observed already durably succeeded and is unaffected"
);
}
#[cfg(test)]
pub(crate) mod tests {
use super::{AppendNotification, EventLogObserver, MutationKind, MutationNotification};
use crate::{EventLogHost, PartitionMigration, RewriteDecision};
use polyc_eventlog::Event;
use std::sync::{Arc, Mutex};
use tokio_util::sync::CancellationToken;
type RecordedAppend = (String, Vec<Event>, Vec<u64>);
type RecordedMutation = (String, MutationKind, u64);
#[derive(Default)]
pub(crate) struct RecordingObserver {
pub(crate) appends: Mutex<Vec<RecordedAppend>>,
pub(crate) mutations: Mutex<Vec<RecordedMutation>>,
}
impl EventLogObserver for RecordingObserver {
fn on_append(&self, notification: &AppendNotification<'_>) {
self.appends.lock().expect("poison").push((
notification.partition.to_owned(),
notification.events.to_vec(),
notification.positions.to_vec(),
));
}
fn on_mutation(&self, notification: &MutationNotification<'_>) {
self.mutations.lock().expect("poison").push((
notification.partition.to_owned(),
notification.kind.clone(),
notification.epoch,
));
}
}
struct PanickingObserver;
impl EventLogObserver for PanickingObserver {
fn on_append(&self, _notification: &AppendNotification<'_>) {
panic!("deliberate test panic: proving the host contains an observer panic");
}
}
#[tokio::test]
async fn a_committed_batch_counts_its_attestation() {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-attestation-count-{}",
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 before = crate::metrics::attestation_count(crate::metrics::AttestationOutcome::Signed);
host.append_batch(
"conv-attestation-count".to_owned(),
vec![Event::new("output_msg", b"one".to_vec())],
)
.await
.expect("append");
assert_eq!(
crate::metrics::attestation_count(crate::metrics::AttestationOutcome::Signed),
before + 1,
"a committed batch must count its signed-root outcome"
);
shutdown.cancel();
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn append_observer_sees_correct_positions_for_single_and_batch_appends() {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-observer-append-{}",
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 recorder = Arc::new(RecordingObserver::default());
host.register_observer(recorder.clone());
let first_positions = host
.append_batch(
"conv-observer-append".to_owned(),
vec![Event::new("output_msg", b"one".to_vec())],
)
.await
.expect("single append");
let batch_positions = host
.append_batch(
"conv-observer-append".to_owned(),
vec![
Event::new("output_msg", b"two".to_vec()),
Event::new("output_msg", b"three".to_vec()),
],
)
.await
.expect("batch append");
let seen = recorder.appends.lock().expect("poison");
assert_eq!(seen.len(), 2, "one notification per append_batch call");
let (partition, events, positions) = &seen[0];
assert_eq!(partition, "conv-observer-append");
assert_eq!(
events.len(),
1,
"the single-event call carries exactly one event"
);
assert_eq!(positions, &first_positions);
let (partition, events, positions) = &seen[1];
assert_eq!(partition, "conv-observer-append");
assert_eq!(positions, &batch_positions);
assert_eq!(
events.iter().map(|e| e.payload.clone()).collect::<Vec<_>>(),
vec![b"two".to_vec(), b"three".to_vec()],
"batch events arrive in append order"
);
drop(seen);
drop(host);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn destroy_and_rewrite_bump_the_epoch_and_notify_append_does_not() {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-observer-mutation-{}",
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 recorder = Arc::new(RecordingObserver::default());
host.register_observer(recorder.clone());
assert_eq!(
host.mutation_epoch("conv-observer-mutation"),
0,
"a partition this process never mutated starts at epoch 0"
);
host.append_batch(
"conv-observer-mutation".to_owned(),
vec![Event::new("memory_added", b"fact".to_vec())],
)
.await
.expect("seed");
assert_eq!(
host.mutation_epoch("conv-observer-mutation"),
0,
"append_batch must never bump the mutation epoch"
);
assert!(
recorder.mutations.lock().expect("poison").is_empty(),
"append_batch must never fire a mutation notification"
);
host.rewrite_partition_for_test(
"conv-observer-mutation".to_owned(),
Box::new(|event| {
if event.kind == "memory_added" {
RewriteDecision::Drop
} else {
RewriteDecision::Keep
}
}),
)
.await
.expect("rewrite");
assert_eq!(
host.mutation_epoch("conv-observer-mutation"),
1,
"rewrite is the first mutation"
);
host.destroy_partition("conv-observer-mutation".to_owned())
.await
.expect("destroy");
assert_eq!(
host.mutation_epoch("conv-observer-mutation"),
2,
"destroy is the second mutation"
);
let mutations = recorder.mutations.lock().expect("poison");
assert_eq!(
mutations.len(),
2,
"one notification for the rewrite, one for the destroy"
);
assert_eq!(mutations[0].0, "conv-observer-mutation");
assert_eq!(mutations[0].1, MutationKind::Rewritten { dropped: 1 });
assert_eq!(mutations[0].2, 1);
assert_eq!(mutations[1].0, "conv-observer-mutation");
assert_eq!(mutations[1].1, MutationKind::Destroyed);
assert_eq!(mutations[1].2, 2);
drop(mutations);
drop(host);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn multiple_observers_all_fire_for_the_same_append() {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-observer-multi-{}",
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 first = Arc::new(RecordingObserver::default());
let second = Arc::new(RecordingObserver::default());
host.register_observer(first.clone());
host.register_observer(second.clone());
host.append_batch(
"conv-observer-multi".to_owned(),
vec![Event::new("output_msg", b"hi".to_vec())],
)
.await
.expect("append");
assert_eq!(
first.appends.lock().expect("poison").len(),
1,
"the first registered observer sees the append"
);
assert_eq!(
second.appends.lock().expect("poison").len(),
1,
"the second registered observer sees the SAME append"
);
drop(host);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn a_panicking_observer_never_blocks_or_poisons_the_write_path() {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-observer-panic-{}",
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 recorder = Arc::new(RecordingObserver::default());
host.register_observer(Arc::new(PanickingObserver));
host.register_observer(recorder.clone());
for i in 0..3u32 {
host.append_batch(
"conv-observer-panic".to_owned(),
vec![Event::new(format!("k{i}"), b"x".to_vec())],
)
.await
.expect("append must succeed even though an observer panics on every call");
}
assert_eq!(
recorder.appends.lock().expect("poison").len(),
3,
"the well-behaved observer keeps seeing every append despite the other one \
panicking every time"
);
drop(host);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn migrate_partition_notifies_both_the_source_and_the_destination() {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-observer-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 recorder = Arc::new(RecordingObserver::default());
host.register_observer(recorder.clone());
host.append_batch(
"persona-observer-migrate-source".to_owned(),
vec![Event::new("memory_added", b"fact-a".to_vec())],
)
.await
.expect("seed source");
let migration = host
.migrate_partition(
"persona-observer-migrate-source".to_owned(),
"persona-observer-migrate-dest".to_owned(),
)
.await
.expect("migrate");
assert_eq!(migration, PartitionMigration::Migrated(1));
assert_eq!(
host.mutation_epoch("persona-observer-migrate-source"),
1,
"the source's destroy bumps its own epoch"
);
assert_eq!(
host.mutation_epoch("persona-observer-migrate-dest"),
1,
"the destination's append-in bumps its own, independent epoch"
);
let mutations = recorder.mutations.lock().expect("poison");
assert_eq!(
mutations.len(),
2,
"one notification per side of the migration"
);
let source_notification = mutations
.iter()
.find(|(partition, ..)| partition == "persona-observer-migrate-source")
.expect("the source is notified");
assert_eq!(
source_notification.1,
MutationKind::MigrationSource {
to: "persona-observer-migrate-dest".to_owned(),
}
);
let destination_notification = mutations
.iter()
.find(|(partition, ..)| partition == "persona-observer-migrate-dest")
.expect("the destination is notified");
assert_eq!(
destination_notification.1,
MutationKind::MigrationDestination {
from: "persona-observer-migrate-source".to_owned(),
appended: 1,
}
);
drop(mutations);
drop(host);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn repair_bumps_the_epoch_and_notifies_when_it_quarantines_something() {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-observer-repair-{}",
std::process::id()
));
let _ = std::fs::remove_dir_all(&dir);
let partition = "conv-observer-repair".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");
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 recorder = Arc::new(RecordingObserver::default());
host2.register_observer(recorder.clone());
assert_eq!(host2.mutation_epoch(&partition), 0);
let quarantined = host2
.repair_partition(partition.clone())
.await
.expect("repair completes");
assert!(
!quarantined.is_empty(),
"the corrupted section must quarantine at least one item for this test to be \
meaningful"
);
assert_eq!(
host2.mutation_epoch(&partition),
1,
"a repair that quarantines something bumps the epoch exactly like a rewrite"
);
let mutations = recorder.mutations.lock().expect("poison");
assert_eq!(mutations.len(), 1, "one notification for the repair");
assert_eq!(mutations[0].0, partition);
assert_eq!(
mutations[0].1,
MutationKind::Repaired {
quarantined: quarantined.len(),
}
);
assert_eq!(mutations[0].2, 1);
drop(mutations);
drop(host2);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn a_notification_names_the_logical_partition_not_its_storage_spelling() {
let dir = std::env::temp_dir().join(format!(
"polychrome-eventlog-observer-logical-{}-{}",
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.clone(),
polyc_crypto::signing_role::JournalAttestationSigner::from_seed(1),
)
.expect("spawn host");
let recorder = Arc::new(RecordingObserver::default());
host.register_observer(recorder.clone());
let partition = "conv-app:persona-1:standup";
host.append_batch(partition.to_owned(), vec![Event::new("k", b"v".to_vec())])
.await
.expect("append");
host.rewrite_partition_for_test(partition.to_owned(), Box::new(|_| RewriteDecision::Keep))
.await
.expect("rewrite");
host.migrate_partition(partition.to_owned(), "conv-app:persona-1:moved".to_owned())
.await
.expect("migrate");
let appends = recorder.appends.lock().expect("poison");
assert_eq!(appends.len(), 1, "one append notification");
assert_eq!(
appends[0].0, partition,
"the append notification names the logical partition"
);
drop(appends);
let mutations = recorder.mutations.lock().expect("poison");
assert!(
mutations.len() >= 2,
"the rewrite and the migration both notify: {mutations:?}"
);
for (name, kind, _) in mutations.iter() {
assert!(
name.contains(':'),
"every notification names the logical partition, and this one does not: {name} for {kind:?}"
);
}
assert!(
mutations.iter().any(|(name, kind, _)| name == partition
&& matches!(kind, MutationKind::Rewritten { .. })),
"the rewrite names its own partition logically: {mutations:?}"
);
assert!(
mutations.iter().any(|(name, kind, _)| name == partition
&& matches!(kind, MutationKind::MigrationSource { .. })),
"and so does the migration source: {mutations:?}"
);
drop(mutations);
assert!(
host.mutation_epoch(partition) >= 2,
"the epoch answers under the LOGICAL name, and every mutation bumped it: {}",
host.mutation_epoch(partition)
);
drop(host);
let _ = std::fs::remove_dir_all(&dir);
}
}