use event_emitter_rs::EventEmitter;
use crate::entity::{Entity, EventRecordError, LocalEvent};
use crate::SourcedResult;
pub struct EntityEmitter {
entity: Entity,
event_emitter: EventEmitter,
events_to_emit: Vec<LocalEvent>,
}
impl Default for EntityEmitter {
fn default() -> Self {
Self::new(Entity::new())
}
}
impl EntityEmitter {
pub fn new(entity: Entity) -> Self {
Self {
entity,
event_emitter: EventEmitter::new(),
events_to_emit: Vec::new(),
}
}
pub fn entity(&self) -> &Entity {
&self.entity
}
pub fn entity_mut(&mut self) -> &mut Entity {
&mut self.entity
}
pub fn into_entity(self) -> Entity {
self.entity
}
pub fn enqueue(&mut self, event_type: impl Into<String>, data: impl Into<String>) {
if self.entity.is_replaying() {
return;
}
self.events_to_emit.push(LocalEvent {
event_type: event_type.into(),
data: data.into(),
});
}
pub fn enqueue_with<T: serde::Serialize>(
&mut self,
event_type: impl Into<String>,
payload: &T,
) -> SourcedResult {
if self.entity.is_replaying() {
return Ok(());
}
let data = serde_json::to_string(payload).map_err(|err| EventRecordError {
message: format!("failed to serialize local event payload to JSON: {}", err),
})?;
self.events_to_emit.push(LocalEvent {
event_type: event_type.into(),
data,
});
Ok(())
}
pub fn drain_queued_events(&mut self) -> Vec<LocalEvent> {
self.events_to_emit.drain(..).collect()
}
pub fn is_replaying(&self) -> bool {
self.entity.is_replaying()
}
pub fn set_replaying(&mut self, replaying: bool) {
self.entity.set_replaying(replaying);
}
pub fn on<F>(&mut self, event: &str, listener: F)
where
F: Fn(String) + Send + Sync + 'static,
{
self.event_emitter.on(event, listener);
}
pub fn emit(&mut self, event: &str, data: impl Into<String>) {
self.event_emitter.emit(event, data.into());
}
pub fn emit_queued(&mut self) {
let events: Vec<_> = self.events_to_emit.drain(..).collect();
for event in events {
self.emit(&event.event_type, event.data);
}
}
pub fn queued_len(&self) -> usize {
self.events_to_emit.len()
}
}
pub trait EmittableEntity {
fn with_emitter(self) -> EntityEmitter;
}
impl EmittableEntity for Entity {
fn with_emitter(self) -> EntityEmitter {
EntityEmitter::new(self)
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use std::thread;
use std::time::Duration;
#[test]
fn enqueue_and_emit() {
let entity = Entity::with_id("test");
let mut emitter = entity.with_emitter();
let called = Arc::new(AtomicBool::new(false));
let flag = Arc::clone(&called);
emitter.on("TestEvent", move |data| {
assert_eq!(data, "test payload");
flag.store(true, Ordering::SeqCst);
});
emitter.enqueue("TestEvent", "test payload");
assert_eq!(emitter.queued_len(), 1);
emitter.emit_queued();
assert_eq!(emitter.queued_len(), 0);
thread::sleep(Duration::from_millis(50));
assert!(called.load(Ordering::SeqCst));
}
#[test]
fn entity_access() {
let entity = Entity::with_id("test");
let mut emitter = entity.with_emitter();
assert_eq!(emitter.entity().id(), "test");
emitter.entity_mut().set_id("changed");
assert_eq!(emitter.entity().id(), "changed");
let entity = emitter.into_entity();
assert_eq!(entity.id(), "changed");
}
#[derive(Debug)]
struct FailingSerialize;
impl serde::Serialize for FailingSerialize {
fn serialize<S>(&self, _serializer: S) -> Result<S::Ok, S::Error>
where
S: serde::Serializer,
{
Err(serde::ser::Error::custom("injected serialize failure"))
}
}
#[test]
fn enqueue_with_returns_serialization_error() {
let entity = Entity::with_id("test");
let mut emitter = entity.with_emitter();
let err = emitter
.enqueue_with("BadEvent", &FailingSerialize)
.unwrap_err();
assert!(err
.message
.contains("failed to serialize local event payload to JSON"));
assert!(err.message.contains("injected serialize failure"));
assert_eq!(emitter.queued_len(), 0);
}
#[test]
fn enqueue_with_does_not_serialize_during_replay() {
let mut entity = Entity::with_id("test");
entity.set_replaying(true);
let mut emitter = entity.with_emitter();
emitter.enqueue_with("BadEvent", &FailingSerialize).unwrap();
assert_eq!(emitter.queued_len(), 0);
}
}