use crate::hooks::HookContext;
use crate::Value;
use std::collections::HashMap;
use std::sync::{Arc, RwLock};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum Event {
BeforeInsert,
AfterInsert,
BeforeUpdate,
AfterUpdate,
BeforeDelete,
AfterDelete,
AfterFind,
BeforeRestore,
AfterRestore,
}
impl Event {
pub fn is_before(&self) -> bool {
matches!(
self,
Event::BeforeInsert | Event::BeforeUpdate | Event::BeforeDelete | Event::BeforeRestore
)
}
pub fn is_after(&self) -> bool {
matches!(
self,
Event::AfterInsert
| Event::AfterUpdate
| Event::AfterDelete
| Event::AfterFind
| Event::AfterRestore
)
}
pub fn is_write_event(&self) -> bool {
matches!(
self,
Event::BeforeInsert
| Event::AfterInsert
| Event::BeforeUpdate
| Event::AfterUpdate
| Event::BeforeDelete
| Event::AfterDelete
)
}
pub fn name(&self) -> &'static str {
match self {
Event::BeforeInsert => "before_insert",
Event::AfterInsert => "after_insert",
Event::BeforeUpdate => "before_update",
Event::AfterUpdate => "after_update",
Event::BeforeDelete => "before_delete",
Event::AfterDelete => "after_delete",
Event::AfterFind => "after_find",
Event::BeforeRestore => "before_restore",
Event::AfterRestore => "after_restore",
}
}
}
#[derive(Debug)]
pub enum SubscriberError {
Failed {
subscriber: String,
reason: String,
},
Vetoed {
subscriber: String,
reason: String,
},
}
impl std::fmt::Display for SubscriberError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
SubscriberError::Failed { subscriber, reason } => {
write!(f, "Subscriber `{}` failed: {}", subscriber, reason)
}
SubscriberError::Vetoed { subscriber, reason } => {
write!(f, "Subscriber `{}` vetoed: {}", subscriber, reason)
}
}
}
}
impl std::error::Error for SubscriberError {}
pub type SubscriberResult<T> = Result<T, SubscriberError>;
pub trait Observer: Send + Sync {
fn name(&self) -> &str {
"anonymous_observer"
}
fn before_insert(
&self,
_ctx: &HookContext,
_attrs: &mut HashMap<String, Value>,
) -> SubscriberResult<()> {
Ok(())
}
fn after_insert(
&self,
_ctx: &HookContext,
_attrs: &HashMap<String, Value>,
) -> SubscriberResult<()> {
Ok(())
}
fn before_update(
&self,
_ctx: &HookContext,
_attrs: &mut HashMap<String, Value>,
) -> SubscriberResult<()> {
Ok(())
}
fn after_update(
&self,
_ctx: &HookContext,
_attrs: &HashMap<String, Value>,
) -> SubscriberResult<()> {
Ok(())
}
fn before_delete(
&self,
_ctx: &HookContext,
_attrs: &HashMap<String, Value>,
) -> SubscriberResult<()> {
Ok(())
}
fn after_delete(
&self,
_ctx: &HookContext,
_attrs: &HashMap<String, Value>,
) -> SubscriberResult<()> {
Ok(())
}
fn after_find(
&self,
_ctx: &HookContext,
_attrs: &mut HashMap<String, Value>,
) -> SubscriberResult<()> {
Ok(())
}
}
pub trait EventSubscriber: Send + Sync {
fn name(&self) -> &str {
"anonymous_subscriber"
}
fn subscribed_events(&self) -> Vec<Event>;
fn on_event(
&self,
event: Event,
ctx: &HookContext,
attrs: &HashMap<String, Value>,
) -> SubscriberResult<()>;
}
pub struct EventDispatcher {
observers: RwLock<Vec<Arc<dyn Observer>>>,
subscribers: RwLock<Vec<Arc<dyn EventSubscriber>>>,
errors: RwLock<Vec<SubscriberError>>,
max_errors: usize,
}
const DEFAULT_MAX_ERRORS: usize = 1024;
impl EventDispatcher {
pub fn new() -> Self {
Self {
observers: RwLock::new(Vec::new()),
subscribers: RwLock::new(Vec::new()),
errors: RwLock::new(Vec::new()),
max_errors: DEFAULT_MAX_ERRORS,
}
}
pub fn with_max_errors(mut self, max_errors: usize) -> Self {
self.max_errors = max_errors;
self
}
pub fn add_observer(&self, observer: Box<dyn Observer>) {
let arc: Arc<dyn Observer> = Arc::from(observer);
self.observers.write().unwrap().push(arc);
}
pub fn subscribe(&self, subscriber: Box<dyn EventSubscriber>) {
let arc: Arc<dyn EventSubscriber> = Arc::from(subscriber);
self.subscribers.write().unwrap().push(arc);
}
pub fn clear(&self) {
self.observers.write().unwrap().clear();
self.subscribers.write().unwrap().clear();
self.errors.write().unwrap().clear();
}
pub fn observer_count(&self) -> usize {
self.observers.read().unwrap().len()
}
pub fn subscriber_count(&self) -> usize {
self.subscribers.read().unwrap().len()
}
pub fn drain_errors(&self) -> Vec<SubscriberError> {
std::mem::take(&mut *self.errors.write().unwrap())
}
pub fn error_count(&self) -> usize {
self.errors.read().unwrap().len()
}
fn push_errors(&self, new_errors: Vec<SubscriberError>) {
if new_errors.is_empty() {
return;
}
let mut errors = self.errors.write().unwrap();
if self.max_errors == 0 {
errors.extend(new_errors);
return;
}
for e in new_errors {
if errors.len() >= self.max_errors {
errors.remove(0);
}
errors.push(e);
}
}
pub fn dispatch(&self, event: Event, ctx: &HookContext, attrs: &HashMap<String, Value>) {
let mut local_errors: Vec<SubscriberError> = Vec::new();
let observers_snapshot: Vec<Arc<dyn Observer>> = {
let observers = self.observers.read().unwrap();
observers.clone()
};
for observer in observers_snapshot.iter() {
let result = match event {
Event::AfterInsert => observer.after_insert(ctx, attrs),
Event::AfterUpdate => observer.after_update(ctx, attrs),
Event::AfterDelete => observer.after_delete(ctx, attrs),
_ => Ok(()),
};
if let Err(e) = result {
local_errors.push(e);
}
}
let subscribers_snapshot: Vec<Arc<dyn EventSubscriber>> = {
let subscribers = self.subscribers.read().unwrap();
subscribers.clone()
};
for subscriber in subscribers_snapshot.iter() {
if !subscriber.subscribed_events().contains(&event) {
continue;
}
if let Err(e) = subscriber.on_event(event, ctx, attrs) {
local_errors.push(e);
}
}
if !local_errors.is_empty() {
self.push_errors(local_errors);
}
}
pub fn dispatch_before_mut(
&self,
event: Event,
ctx: &HookContext,
attrs: &mut HashMap<String, Value>,
) -> SubscriberResult<()> {
let mut local_errors: Vec<SubscriberError> = Vec::new();
let mut vetoed: Option<SubscriberError> = None;
let observers_snapshot: Vec<Arc<dyn Observer>> = {
let observers = self.observers.read().unwrap();
observers.clone()
};
for observer in observers_snapshot.iter() {
let result = match event {
Event::BeforeInsert => observer.before_insert(ctx, attrs),
Event::BeforeUpdate => observer.before_update(ctx, attrs),
_ => Ok(()),
};
match result {
Ok(()) => {}
Err(e @ SubscriberError::Vetoed { .. }) => {
vetoed = Some(e);
break;
}
Err(e) => local_errors.push(e),
}
}
if vetoed.is_none() {
let subscribers_snapshot: Vec<Arc<dyn EventSubscriber>> = {
let subscribers = self.subscribers.read().unwrap();
subscribers.clone()
};
for subscriber in subscribers_snapshot.iter() {
if !subscriber.subscribed_events().contains(&event) {
continue;
}
match subscriber.on_event(event, ctx, attrs) {
Ok(()) => {}
Err(e @ SubscriberError::Vetoed { .. }) => {
vetoed = Some(e);
break;
}
Err(e) => local_errors.push(e),
}
}
}
if !local_errors.is_empty() {
self.push_errors(local_errors);
}
if let Some(e) = vetoed {
return Err(e);
}
Ok(())
}
pub fn dispatch_after_find(
&self,
ctx: &HookContext,
attrs: &mut HashMap<String, Value>,
) -> SubscriberResult<()> {
let mut local_errors: Vec<SubscriberError> = Vec::new();
let observers_snapshot: Vec<Arc<dyn Observer>> = {
let observers = self.observers.read().unwrap();
observers.clone()
};
for observer in observers_snapshot.iter() {
if let Err(e) = observer.after_find(ctx, attrs) {
local_errors.push(e);
}
}
let subscribers_snapshot: Vec<Arc<dyn EventSubscriber>> = {
let subscribers = self.subscribers.read().unwrap();
subscribers.clone()
};
for subscriber in subscribers_snapshot.iter() {
if !subscriber.subscribed_events().contains(&Event::AfterFind) {
continue;
}
if let Err(e) = subscriber.on_event(Event::AfterFind, ctx, attrs) {
local_errors.push(e);
}
}
if !local_errors.is_empty() {
self.push_errors(local_errors);
}
Ok(())
}
}
impl Default for EventDispatcher {
fn default() -> Self {
Self::new()
}
}
#[derive(Clone)]
pub struct AuditLogSubscriber {
logs: std::sync::Arc<std::sync::Mutex<Vec<String>>>,
}
impl AuditLogSubscriber {
pub fn new() -> Self {
Self {
logs: std::sync::Arc::new(std::sync::Mutex::new(Vec::new())),
}
}
pub fn logs(&self) -> &std::sync::Arc<std::sync::Mutex<Vec<String>>> {
&self.logs
}
}
impl Default for AuditLogSubscriber {
fn default() -> Self {
Self::new()
}
}
impl EventSubscriber for AuditLogSubscriber {
fn name(&self) -> &str {
"audit_log"
}
fn subscribed_events(&self) -> Vec<Event> {
vec![Event::AfterInsert, Event::AfterUpdate, Event::AfterDelete]
}
fn on_event(
&self,
event: Event,
ctx: &HookContext,
attrs: &HashMap<String, Value>,
) -> SubscriberResult<()> {
let mut logs = self.logs.lock().unwrap();
logs.push(format!(
"event={} operator={:?} field_count={}",
event.name(),
ctx.operator_id,
attrs.len()
));
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::{Arc, Mutex};
#[test]
fn test_event_is_before_after() {
assert!(Event::BeforeInsert.is_before());
assert!(!Event::BeforeInsert.is_after());
assert!(Event::AfterInsert.is_after());
assert!(!Event::AfterInsert.is_before());
}
#[test]
fn test_event_is_write_event() {
assert!(Event::BeforeInsert.is_write_event());
assert!(Event::AfterUpdate.is_write_event());
assert!(Event::BeforeDelete.is_write_event());
assert!(!Event::AfterFind.is_write_event());
}
#[test]
fn test_event_name() {
assert_eq!(Event::BeforeInsert.name(), "before_insert");
assert_eq!(Event::AfterDelete.name(), "after_delete");
assert_eq!(Event::AfterFind.name(), "after_find");
}
#[test]
fn test_new_dispatcher_is_empty() {
let d = EventDispatcher::new();
assert_eq!(d.observer_count(), 0);
assert_eq!(d.subscriber_count(), 0);
}
#[test]
fn test_add_observer() {
struct DummyObserver;
impl Observer for DummyObserver {}
let d = EventDispatcher::new();
d.add_observer(Box::new(DummyObserver));
assert_eq!(d.observer_count(), 1);
}
#[test]
fn test_subscribe() {
struct DummySubscriber;
impl EventSubscriber for DummySubscriber {
fn subscribed_events(&self) -> Vec<Event> {
vec![Event::AfterInsert]
}
fn on_event(
&self,
_event: Event,
_ctx: &HookContext,
_attrs: &HashMap<String, Value>,
) -> SubscriberResult<()> {
Ok(())
}
}
let d = EventDispatcher::new();
d.subscribe(Box::new(DummySubscriber));
assert_eq!(d.subscriber_count(), 1);
}
#[test]
fn test_clear() {
struct DummyObserver;
impl Observer for DummyObserver {}
let d = EventDispatcher::new();
d.add_observer(Box::new(DummyObserver));
d.clear();
assert_eq!(d.observer_count(), 0);
}
struct CountingObserver {
insert_count: Arc<Mutex<u32>>,
update_count: Arc<Mutex<u32>>,
delete_count: Arc<Mutex<u32>>,
}
impl Observer for CountingObserver {
fn name(&self) -> &str {
"counting"
}
fn after_insert(
&self,
_ctx: &HookContext,
_attrs: &HashMap<String, Value>,
) -> SubscriberResult<()> {
*self.insert_count.lock().unwrap() += 1;
Ok(())
}
fn after_update(
&self,
_ctx: &HookContext,
_attrs: &HashMap<String, Value>,
) -> SubscriberResult<()> {
*self.update_count.lock().unwrap() += 1;
Ok(())
}
fn after_delete(
&self,
_ctx: &HookContext,
_attrs: &HashMap<String, Value>,
) -> SubscriberResult<()> {
*self.delete_count.lock().unwrap() += 1;
Ok(())
}
}
#[test]
fn test_observer_triggered_on_dispatch() {
let insert = Arc::new(Mutex::new(0u32));
let update = Arc::new(Mutex::new(0u32));
let delete = Arc::new(Mutex::new(0u32));
let observer = CountingObserver {
insert_count: insert.clone(),
update_count: update.clone(),
delete_count: delete.clone(),
};
let d = EventDispatcher::new();
d.add_observer(Box::new(observer));
let ctx = HookContext::default();
let attrs = HashMap::new();
d.dispatch(Event::AfterInsert, &ctx, &attrs);
d.dispatch(Event::AfterInsert, &ctx, &attrs);
d.dispatch(Event::AfterUpdate, &ctx, &attrs);
d.dispatch(Event::AfterDelete, &ctx, &attrs);
assert_eq!(*insert.lock().unwrap(), 2);
assert_eq!(*update.lock().unwrap(), 1);
assert_eq!(*delete.lock().unwrap(), 1);
}
#[test]
fn test_observer_before_event_can_modify_attrs() {
struct TimestampInjector;
impl Observer for TimestampInjector {
fn before_insert(
&self,
_ctx: &HookContext,
attrs: &mut HashMap<String, Value>,
) -> SubscriberResult<()> {
attrs.insert(
"created_at".to_string(),
Value::String("2026-07-19".to_string()),
);
Ok(())
}
}
let d = EventDispatcher::new();
d.add_observer(Box::new(TimestampInjector));
let ctx = HookContext::default();
let mut attrs = HashMap::new();
d.dispatch_before_mut(Event::BeforeInsert, &ctx, &mut attrs)
.unwrap();
assert_eq!(
attrs.get("created_at"),
Some(&Value::String("2026-07-19".to_string()))
);
}
struct InsertOnlySubscriber {
called: Arc<Mutex<u32>>,
}
impl EventSubscriber for InsertOnlySubscriber {
fn name(&self) -> &str {
"insert_only"
}
fn subscribed_events(&self) -> Vec<Event> {
vec![Event::AfterInsert]
}
fn on_event(
&self,
_event: Event,
_ctx: &HookContext,
_attrs: &HashMap<String, Value>,
) -> SubscriberResult<()> {
*self.called.lock().unwrap() += 1;
Ok(())
}
}
#[test]
fn test_subscriber_only_called_for_subscribed_events() {
let called = Arc::new(Mutex::new(0u32));
let subscriber = InsertOnlySubscriber {
called: called.clone(),
};
let d = EventDispatcher::new();
d.subscribe(Box::new(subscriber));
let ctx = HookContext::default();
let attrs = HashMap::new();
d.dispatch(Event::AfterInsert, &ctx, &attrs);
d.dispatch(Event::AfterUpdate, &ctx, &attrs);
d.dispatch(Event::AfterDelete, &ctx, &attrs);
d.dispatch(Event::AfterInsert, &ctx, &attrs);
assert_eq!(*called.lock().unwrap(), 2);
}
#[test]
fn test_subscriber_veto_aborts_before_event() {
struct VetoSubscriber;
impl EventSubscriber for VetoSubscriber {
fn name(&self) -> &str {
"veto"
}
fn subscribed_events(&self) -> Vec<Event> {
vec![Event::BeforeInsert]
}
fn on_event(
&self,
_event: Event,
_ctx: &HookContext,
_attrs: &HashMap<String, Value>,
) -> SubscriberResult<()> {
Err(SubscriberError::Vetoed {
subscriber: "veto".to_string(),
reason: "Business rule violation".to_string(),
})
}
}
let d = EventDispatcher::new();
d.subscribe(Box::new(VetoSubscriber));
let ctx = HookContext::default();
let mut attrs = HashMap::new();
let result = d.dispatch_before_mut(Event::BeforeInsert, &ctx, &mut attrs);
assert!(matches!(result, Err(SubscriberError::Vetoed { .. })));
}
#[test]
fn test_subscriber_failed_does_not_abort_after_event() {
struct FailingSubscriber;
impl EventSubscriber for FailingSubscriber {
fn name(&self) -> &str {
"failing"
}
fn subscribed_events(&self) -> Vec<Event> {
vec![Event::AfterInsert]
}
fn on_event(
&self,
_event: Event,
_ctx: &HookContext,
_attrs: &HashMap<String, Value>,
) -> SubscriberResult<()> {
Err(SubscriberError::Failed {
subscriber: "failing".to_string(),
reason: "Connection lost".to_string(),
})
}
}
struct CountingSubscriber {
called: Arc<Mutex<u32>>,
}
impl EventSubscriber for CountingSubscriber {
fn name(&self) -> &str {
"counting"
}
fn subscribed_events(&self) -> Vec<Event> {
vec![Event::AfterInsert]
}
fn on_event(
&self,
_event: Event,
_ctx: &HookContext,
_attrs: &HashMap<String, Value>,
) -> SubscriberResult<()> {
*self.called.lock().unwrap() += 1;
Ok(())
}
}
let called = Arc::new(Mutex::new(0u32));
let d = EventDispatcher::new();
d.subscribe(Box::new(FailingSubscriber));
d.subscribe(Box::new(CountingSubscriber {
called: called.clone(),
}));
let ctx = HookContext::default();
let attrs = HashMap::new();
d.dispatch(Event::AfterInsert, &ctx, &attrs);
assert_eq!(*called.lock().unwrap(), 1);
}
#[test]
fn test_audit_log_subscriber() {
let audit = AuditLogSubscriber::new();
let audit_clone = audit.clone();
let d = EventDispatcher::new();
d.subscribe(Box::new(audit_clone));
let ctx = HookContext {
operator_id: Some(42),
..Default::default()
};
let mut attrs = HashMap::new();
attrs.insert("name".to_string(), Value::String("alice".to_string()));
d.dispatch(Event::AfterInsert, &ctx, &attrs);
d.dispatch(Event::AfterUpdate, &ctx, &attrs);
d.dispatch(Event::AfterDelete, &ctx, &attrs);
d.dispatch_after_find(&ctx, &mut attrs).unwrap();
let logs = audit.logs().lock().unwrap();
assert_eq!(logs.len(), 3);
assert!(logs[0].contains("event=after_insert"));
assert!(logs[0].contains("operator=Some(42)"));
assert!(logs[0].contains("field_count=1"));
}
#[test]
fn test_multiple_subscribers_and_observers() {
let sub1_called = Arc::new(Mutex::new(0u32));
let sub2_called = Arc::new(Mutex::new(0u32));
let obs_called = Arc::new(Mutex::new(0u32));
struct Sub1(Arc<Mutex<u32>>);
impl EventSubscriber for Sub1 {
fn name(&self) -> &str {
"sub1"
}
fn subscribed_events(&self) -> Vec<Event> {
vec![Event::AfterInsert]
}
fn on_event(
&self,
_e: Event,
_c: &HookContext,
_a: &HashMap<String, Value>,
) -> SubscriberResult<()> {
*self.0.lock().unwrap() += 1;
Ok(())
}
}
struct Sub2(Arc<Mutex<u32>>);
impl EventSubscriber for Sub2 {
fn name(&self) -> &str {
"sub2"
}
fn subscribed_events(&self) -> Vec<Event> {
vec![Event::AfterInsert, Event::AfterUpdate]
}
fn on_event(
&self,
_e: Event,
_c: &HookContext,
_a: &HashMap<String, Value>,
) -> SubscriberResult<()> {
*self.0.lock().unwrap() += 1;
Ok(())
}
}
struct Obs(Arc<Mutex<u32>>);
impl Observer for Obs {
fn name(&self) -> &str {
"obs"
}
fn after_insert(
&self,
_c: &HookContext,
_a: &HashMap<String, Value>,
) -> SubscriberResult<()> {
*self.0.lock().unwrap() += 1;
Ok(())
}
}
let d = EventDispatcher::new();
d.subscribe(Box::new(Sub1(sub1_called.clone())));
d.subscribe(Box::new(Sub2(sub2_called.clone())));
d.add_observer(Box::new(Obs(obs_called.clone())));
let ctx = HookContext::default();
let attrs = HashMap::new();
d.dispatch(Event::AfterInsert, &ctx, &attrs);
assert_eq!(*sub1_called.lock().unwrap(), 1);
assert_eq!(*sub2_called.lock().unwrap(), 1);
assert_eq!(*obs_called.lock().unwrap(), 1);
}
#[test]
fn test_drain_errors() {
struct ErrSub;
impl EventSubscriber for ErrSub {
fn name(&self) -> &str {
"err"
}
fn subscribed_events(&self) -> Vec<Event> {
vec![Event::AfterInsert]
}
fn on_event(
&self,
_e: Event,
_c: &HookContext,
_a: &HashMap<String, Value>,
) -> SubscriberResult<()> {
Err(SubscriberError::Failed {
subscriber: "err".to_string(),
reason: "test".to_string(),
})
}
}
let d = EventDispatcher::new();
d.subscribe(Box::new(ErrSub));
let ctx = HookContext::default();
let attrs = HashMap::new();
d.dispatch(Event::AfterInsert, &ctx, &attrs);
d.dispatch(Event::AfterInsert, &ctx, &attrs);
let errors = d.drain_errors();
assert_eq!(errors.len(), 2);
assert!(matches!(errors[0], SubscriberError::Failed { .. }));
let errors = d.drain_errors();
assert!(errors.is_empty());
}
#[test]
fn test_max_errors_limits_buffer_size() {
struct ErrSub;
impl EventSubscriber for ErrSub {
fn name(&self) -> &str {
"err"
}
fn subscribed_events(&self) -> Vec<Event> {
vec![Event::AfterInsert]
}
fn on_event(
&self,
_e: Event,
_c: &HookContext,
_a: &HashMap<String, Value>,
) -> SubscriberResult<()> {
Err(SubscriberError::Failed {
subscriber: "err".to_string(),
reason: "test".to_string(),
})
}
}
let d = EventDispatcher::new().with_max_errors(3);
d.subscribe(Box::new(ErrSub));
let ctx = HookContext::default();
let attrs = HashMap::new();
for _ in 0..5 {
d.dispatch(Event::AfterInsert, &ctx, &attrs);
}
assert_eq!(d.error_count(), 3);
let errors = d.drain_errors();
assert_eq!(errors.len(), 3);
}
#[test]
fn test_max_errors_zero_means_unlimited() {
struct ErrSub;
impl EventSubscriber for ErrSub {
fn name(&self) -> &str {
"err"
}
fn subscribed_events(&self) -> Vec<Event> {
vec![Event::AfterInsert]
}
fn on_event(
&self,
_e: Event,
_c: &HookContext,
_a: &HashMap<String, Value>,
) -> SubscriberResult<()> {
Err(SubscriberError::Failed {
subscriber: "err".to_string(),
reason: "test".to_string(),
})
}
}
let d = EventDispatcher::new().with_max_errors(0);
d.subscribe(Box::new(ErrSub));
let ctx = HookContext::default();
let attrs = HashMap::new();
for _ in 0..10 {
d.dispatch(Event::AfterInsert, &ctx, &attrs);
}
assert_eq!(d.error_count(), 10);
}
#[test]
fn test_max_errors_fifo_eviction_order() {
struct CounterSub(Arc<Mutex<u32>>);
impl EventSubscriber for CounterSub {
fn name(&self) -> &str {
"counter"
}
fn subscribed_events(&self) -> Vec<Event> {
vec![Event::AfterInsert]
}
fn on_event(
&self,
_e: Event,
_c: &HookContext,
_a: &HashMap<String, Value>,
) -> SubscriberResult<()> {
let mut n = self.0.lock().unwrap();
*n += 1;
Err(SubscriberError::Failed {
subscriber: "counter".to_string(),
reason: format!("call-{}", *n),
})
}
}
let counter = Arc::new(Mutex::new(0u32));
let d = EventDispatcher::new().with_max_errors(2);
d.subscribe(Box::new(CounterSub(counter.clone())));
let ctx = HookContext::default();
let attrs = HashMap::new();
for _ in 0..4 {
d.dispatch(Event::AfterInsert, &ctx, &attrs);
}
let errors = d.drain_errors();
assert_eq!(errors.len(), 2);
match &errors[0] {
SubscriberError::Failed { reason, .. } => assert_eq!(reason, "call-3"),
other => panic!("expected Failed, got {:?}", other),
}
match &errors[1] {
SubscriberError::Failed { reason, .. } => assert_eq!(reason, "call-4"),
other => panic!("expected Failed, got {:?}", other),
}
}
#[test]
fn test_veto_aborts_subsequent_observers() {
let second_called = Arc::new(Mutex::new(0u32));
struct VetoObs;
impl Observer for VetoObs {
fn name(&self) -> &str {
"veto"
}
fn before_insert(
&self,
_c: &HookContext,
_a: &mut HashMap<String, Value>,
) -> SubscriberResult<()> {
Err(SubscriberError::Vetoed {
subscriber: "veto".to_string(),
reason: "no".to_string(),
})
}
}
struct CountingObs(Arc<Mutex<u32>>);
impl Observer for CountingObs {
fn name(&self) -> &str {
"counting"
}
fn before_insert(
&self,
_c: &HookContext,
_a: &mut HashMap<String, Value>,
) -> SubscriberResult<()> {
*self.0.lock().unwrap() += 1;
Ok(())
}
}
let d = EventDispatcher::new();
d.add_observer(Box::new(VetoObs));
d.add_observer(Box::new(CountingObs(second_called.clone())));
let ctx = HookContext::default();
let mut attrs = HashMap::new();
let result = d.dispatch_before_mut(Event::BeforeInsert, &ctx, &mut attrs);
assert!(result.is_err());
assert_eq!(*second_called.lock().unwrap(), 0);
}
#[test]
fn test_error_display() {
let e = SubscriberError::Failed {
subscriber: "test".to_string(),
reason: "boom".to_string(),
};
assert!(e.to_string().contains("test"));
assert!(e.to_string().contains("boom"));
let e = SubscriberError::Vetoed {
subscriber: "vetoer".to_string(),
reason: "rejected".to_string(),
};
assert!(e.to_string().contains("vetoer"));
assert!(e.to_string().contains("rejected"));
}
}