use std::sync::Arc;
use helix_core::effect::DomainEventBytes;
use helix_core::ports::EventSink;
use crate::engine::BatchSink;
pub struct RecordingSink<S> {
inner: S,
observe: Arc<dyn Fn(&DomainEventBytes) + Send + Sync>,
}
impl<S> RecordingSink<S> {
pub fn new(inner: S, observe: Arc<dyn Fn(&DomainEventBytes) + Send + Sync>) -> Self {
Self { inner, observe }
}
pub fn into_inner(self) -> S {
self.inner
}
}
impl<S: EventSink> EventSink for RecordingSink<S> {
fn emit(&self, event: DomainEventBytes) {
(self.observe)(&event);
self.inner.emit(event);
}
}
impl<S: BatchSink> BatchSink for RecordingSink<S> {
fn flush(&self) {
self.inner.flush();
}
}
#[cfg(test)]
mod tests {
use super::*;
use bytes::Bytes;
use std::sync::atomic::{AtomicUsize, Ordering};
#[derive(Default)]
struct SpySink {
emits: AtomicUsize,
flushes: AtomicUsize,
}
impl EventSink for SpySink {
fn emit(&self, _event: DomainEventBytes) {
self.emits.fetch_add(1, Ordering::Relaxed);
}
}
impl BatchSink for SpySink {
fn flush(&self) {
self.flushes.fetch_add(1, Ordering::Relaxed);
}
}
#[test]
fn recording_sink_observes_and_forwards_emit() {
let observed = Arc::new(AtomicUsize::new(0));
let observed_clone = Arc::clone(&observed);
let inner = SpySink::default();
let sink = RecordingSink::new(
inner,
Arc::new(move |_ev: &DomainEventBytes| {
observed_clone.fetch_add(1, Ordering::Relaxed);
}),
);
sink.emit(DomainEventBytes(Bytes::from_static(b"a")));
sink.emit(DomainEventBytes(Bytes::from_static(b"b")));
sink.flush();
assert_eq!(
observed.load(Ordering::Relaxed),
2,
"observe 应看到每条 emit"
);
let inner = sink.into_inner();
assert_eq!(
inner.emits.load(Ordering::Relaxed),
2,
"内层应收到全部 emit"
);
assert_eq!(inner.flushes.load(Ordering::Relaxed), 1, "flush 应透传内层");
}
}