type EventHandler<T> =
std::boxed::Box<dyn Fn(crate::ports::events::CloudEventsEnvelope<T>) -> crate::HexResult<()>>;
pub struct InMemoryEventBus<T> {
queue:
std::cell::RefCell<std::collections::VecDeque<crate::ports::events::CloudEventsEnvelope<T>>>,
handlers:
std::cell::RefCell<indexmap::IndexMap<std::string::String, std::vec::Vec<EventHandler<T>>>>,
topic: std::string::String,
}
impl<T> InMemoryEventBus<T> {
pub fn new() -> Self {
Self {
queue: std::cell::RefCell::new(std::collections::VecDeque::new()),
handlers: std::cell::RefCell::new(indexmap::IndexMap::new()),
topic: std::string::String::from("default.events"),
}
}
pub fn with_topic(topic: std::string::String) -> Self {
Self {
queue: std::cell::RefCell::new(std::collections::VecDeque::new()),
handlers: std::cell::RefCell::new(indexmap::IndexMap::new()),
topic,
}
}
pub fn queue_size(&self) -> usize {
self.queue.borrow().len()
}
pub fn clear(&mut self) {
self.queue.borrow_mut().clear();
}
}
impl<T> Default for InMemoryEventBus<T> {
fn default() -> Self {
Self::new()
}
}
impl<T> crate::adapters::Adapter for InMemoryEventBus<T> {}
impl<T> crate::ports::events::EventPublisher<T> for InMemoryEventBus<T>
where
T: Clone,
{
fn publish(
&self,
envelope: &crate::ports::events::CloudEventsEnvelope<T>,
) -> crate::HexResult<()> {
envelope.validate()?;
envelope.validate_time_format()?;
let route = if envelope.r#type.is_empty() {
self.topic.as_str()
} else {
envelope.r#type.as_str()
};
let delivered = {
let handlers = self.handlers.borrow();
match handlers.get(route) {
std::option::Option::Some(topic_handlers) if !topic_handlers.is_empty() => {
for handler in topic_handlers {
handler(envelope.clone())?;
}
true
}
_ => false,
}
};
if !delivered {
self.queue.borrow_mut().push_back(envelope.clone());
}
std::result::Result::Ok(())
}
fn publish_batch(
&self,
envelopes: &[crate::ports::events::CloudEventsEnvelope<T>],
) -> crate::HexResult<()> {
for envelope in envelopes {
self.publish(envelope)?;
}
std::result::Result::Ok(())
}
}
impl<T> crate::ports::events::EventSubscriber<T> for InMemoryEventBus<T>
where
T: Clone,
{
fn subscribe(
&mut self,
topic: &str,
handler: std::boxed::Box<
dyn Fn(crate::ports::events::CloudEventsEnvelope<T>) -> crate::HexResult<()>,
>,
) -> crate::HexResult<()> {
self
.handlers
.borrow_mut()
.entry(std::string::String::from(topic))
.or_default()
.push(handler);
std::result::Result::Ok(())
}
fn poll(
&mut self,
) -> crate::HexResult<std::option::Option<crate::ports::events::CloudEventsEnvelope<T>>> {
std::result::Result::Ok(self.queue.borrow_mut().pop_front())
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::ports::events::{EventPublisher, EventSubscriber};
#[derive(Clone)]
struct TestEvent {
id: std::string::String,
value: std::string::String,
}
impl crate::domain::DomainEvent for TestEvent {
fn event_type(&self) -> &str {
"com.test.event.created"
}
fn aggregate_id(&self) -> std::string::String {
self.id.clone()
}
}
#[test]
fn test_new_bus_is_empty() {
let bus: InMemoryEventBus<TestEvent> = InMemoryEventBus::new();
std::assert_eq!(bus.queue_size(), 0);
}
#[test]
fn test_with_topic_sets_topic() {
let bus: InMemoryEventBus<TestEvent> =
InMemoryEventBus::with_topic(std::string::String::from("test.events"));
std::assert_eq!(bus.topic, "test.events");
}
#[test]
fn test_publish_adds_to_queue() {
let bus: InMemoryEventBus<TestEvent> = InMemoryEventBus::new();
let event = TestEvent {
id: std::string::String::from("test-123"),
value: std::string::String::from("test value"),
};
let envelope = crate::ports::events::CloudEventsEnvelope::from_domain_event(
std::string::String::from("evt-001"),
std::string::String::from("/test/source"),
event,
);
bus.publish(&envelope).unwrap();
std::assert_eq!(bus.queue_size(), 1);
}
#[test]
fn test_publish_batch_adds_multiple() {
let bus: InMemoryEventBus<TestEvent> = InMemoryEventBus::new();
let envelopes = vec![
crate::ports::events::CloudEventsEnvelope::from_domain_event(
std::string::String::from("evt-001"),
std::string::String::from("/test/source"),
TestEvent {
id: std::string::String::from("test-1"),
value: std::string::String::from("value-1"),
},
),
crate::ports::events::CloudEventsEnvelope::from_domain_event(
std::string::String::from("evt-002"),
std::string::String::from("/test/source"),
TestEvent {
id: std::string::String::from("test-2"),
value: std::string::String::from("value-2"),
},
),
];
bus.publish_batch(&envelopes).unwrap();
std::assert_eq!(bus.queue_size(), 2);
}
#[test]
fn test_poll_returns_none_when_empty() {
let mut bus: InMemoryEventBus<TestEvent> = InMemoryEventBus::new();
let result = bus.poll().unwrap();
std::assert!(result.is_none());
}
#[test]
fn test_poll_returns_event_and_removes_from_queue() {
let mut bus: InMemoryEventBus<TestEvent> = InMemoryEventBus::new();
let event = TestEvent {
id: std::string::String::from("test-123"),
value: std::string::String::from("test value"),
};
let envelope = crate::ports::events::CloudEventsEnvelope::from_domain_event(
std::string::String::from("evt-001"),
std::string::String::from("/test/source"),
event,
);
bus.publish(&envelope).unwrap();
std::assert_eq!(bus.queue_size(), 1);
let polled = bus.poll().unwrap();
std::assert!(polled.is_some());
std::assert_eq!(bus.queue_size(), 0);
}
#[test]
fn test_subscribe_and_handler_invoked() {
let mut bus: InMemoryEventBus<TestEvent> = InMemoryEventBus::new();
let invoked = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
let invoked_clone = invoked.clone();
bus
.subscribe(
"com.test.event.created",
std::boxed::Box::new(move |_envelope| {
invoked_clone.store(true, std::sync::atomic::Ordering::SeqCst);
std::result::Result::Ok(())
}),
)
.unwrap();
let event = TestEvent {
id: std::string::String::from("test-123"),
value: std::string::String::from("test value"),
};
let envelope = crate::ports::events::CloudEventsEnvelope::from_domain_event(
std::string::String::from("evt-001"),
std::string::String::from("/test/source"),
event,
);
bus.publish(&envelope).unwrap();
std::assert!(invoked.load(std::sync::atomic::Ordering::SeqCst));
}
fn envelope_typed(id: &str, ty: &str) -> crate::ports::events::CloudEventsEnvelope<TestEvent> {
let mut env = crate::ports::events::CloudEventsEnvelope::from_domain_event(
std::string::String::from(id),
std::string::String::from("/test/source"),
TestEvent {
id: std::string::String::from(id),
value: std::string::String::from("v"),
},
);
env.r#type = std::string::String::from(ty);
env
}
#[test]
fn test_multiple_topics_route_independently() {
let mut bus: InMemoryEventBus<TestEvent> = InMemoryEventBus::new();
let a = std::rc::Rc::new(std::cell::Cell::new(0));
let b = std::rc::Rc::new(std::cell::Cell::new(0));
let a_c = a.clone();
let b_c = b.clone();
bus
.subscribe(
"topic.a",
std::boxed::Box::new(move |_e| {
a_c.set(a_c.get() + 1);
std::result::Result::Ok(())
}),
)
.unwrap();
bus
.subscribe(
"topic.b",
std::boxed::Box::new(move |_e| {
b_c.set(b_c.get() + 1);
std::result::Result::Ok(())
}),
)
.unwrap();
bus.publish(&envelope_typed("1", "topic.a")).unwrap();
bus.publish(&envelope_typed("2", "topic.b")).unwrap();
bus.publish(&envelope_typed("3", "topic.a")).unwrap();
assert_eq!(a.get(), 2, "topic.a handler must fire for its two events");
assert_eq!(b.get(), 1, "topic.b handler must fire for its one event");
assert_eq!(bus.queue_size(), 0);
}
#[test]
fn test_multiple_handlers_same_topic_all_fire() {
let mut bus: InMemoryEventBus<TestEvent> = InMemoryEventBus::new();
let count = std::rc::Rc::new(std::cell::Cell::new(0));
for _ in 0..3 {
let c = count.clone();
bus
.subscribe(
"topic.x",
std::boxed::Box::new(move |_e| {
c.set(c.get() + 1);
std::result::Result::Ok(())
}),
)
.unwrap();
}
bus.publish(&envelope_typed("1", "topic.x")).unwrap();
assert_eq!(count.get(), 3, "all three handlers on topic.x must fire");
}
#[test]
fn test_same_topic_handlers_fire_in_registration_order() {
let mut bus: InMemoryEventBus<TestEvent> = InMemoryEventBus::new();
let order = std::rc::Rc::new(std::cell::RefCell::new(std::vec::Vec::new()));
let order_first = order.clone();
bus
.subscribe(
"topic.ordered",
std::boxed::Box::new(move |_e| {
order_first.borrow_mut().push("first");
std::result::Result::Ok(())
}),
)
.unwrap();
let order_second = order.clone();
bus
.subscribe(
"topic.ordered",
std::boxed::Box::new(move |_e| {
order_second.borrow_mut().push("second");
std::result::Result::Ok(())
}),
)
.unwrap();
bus.publish(&envelope_typed("1", "topic.ordered")).unwrap();
assert_eq!(
*order.borrow(),
std::vec!["first", "second"],
"handlers subscribed to the same topic must fire in registration order"
);
}
#[test]
fn test_undelivered_events_queue_delivered_events_do_not() {
let mut bus: InMemoryEventBus<TestEvent> = InMemoryEventBus::new();
bus
.subscribe(
"topic.handled",
std::boxed::Box::new(|_e| std::result::Result::Ok(())),
)
.unwrap();
bus.publish(&envelope_typed("1", "topic.handled")).unwrap();
bus
.publish(&envelope_typed("2", "topic.unhandled"))
.unwrap();
assert_eq!(bus.queue_size(), 1, "only the unhandled event is queued");
let polled = bus.poll().unwrap().expect("one queued event");
assert_eq!(polled.r#type, "topic.unhandled");
}
#[test]
fn test_clear_empties_queue() {
let mut bus: InMemoryEventBus<TestEvent> = InMemoryEventBus::new();
let event = TestEvent {
id: std::string::String::from("test-123"),
value: std::string::String::from("test value"),
};
let envelope = crate::ports::events::CloudEventsEnvelope::from_domain_event(
std::string::String::from("evt-001"),
std::string::String::from("/test/source"),
event,
);
bus.publish(&envelope).unwrap();
std::assert_eq!(bus.queue_size(), 1);
bus.clear();
std::assert_eq!(bus.queue_size(), 0);
}
#[test]
fn test_publish_validates_envelope() {
let bus: InMemoryEventBus<TestEvent> = InMemoryEventBus::new();
let envelope = crate::ports::events::CloudEventsEnvelope::<TestEvent>::new(
std::string::String::from(""),
std::string::String::from("/test/source"),
std::string::String::from("com.test.event"),
);
let result = bus.publish(&envelope);
std::assert!(result.is_err());
std::assert_eq!(bus.queue_size(), 0);
}
}