use bevy::prelude::*;
use std::collections::HashMap;
use std::sync::Mutex;
use crate::{
BusEvent, EventBusError,
backends::{EventBusBackendResource, EventBusBackendExt},
};
#[derive(Resource, Default)]
pub struct EventBusErrorQueue {
pending_errors: Mutex<Vec<Box<dyn Fn(&mut World) + Send + Sync>>>,
}
impl EventBusErrorQueue {
pub fn add_error<T: BusEvent + Event>(&self, error: EventBusError<T>) {
if let Ok(mut pending) = self.pending_errors.lock() {
pending.push(Box::new(move |world: &mut World| {
world.send_event(error.clone());
}));
}
}
pub fn flush_errors(&self, world: &mut World) {
if let Ok(mut pending) = self.pending_errors.lock() {
for error_fn in pending.drain(..) {
error_fn(world);
}
}
}
pub fn drain_pending(&self) -> Vec<Box<dyn Fn(&mut World) + Send + Sync>> {
if let Ok(mut pending) = self.pending_errors.lock() {
std::mem::take(&mut *pending)
} else {
Vec::new()
}
}
}
#[derive(bevy::ecs::system::SystemParam)]
pub struct EventBusWriter<'w, T: BusEvent + Event> {
backend: Option<Res<'w, EventBusBackendResource>>,
events: EventWriter<'w, T>,
error_queue: Res<'w, EventBusErrorQueue>,
}
impl<'w, T: BusEvent + Event> EventBusWriter<'w, T> {
pub fn write(&mut self, topic: &str, event: T) {
let event_clone = event.clone();
if let Some(backend_res) = &self.backend {
let backend = backend_res.read();
if !backend.try_send(&event, topic) {
let error_event = EventBusError::immediate(
topic.to_string(),
crate::EventBusErrorType::Other, "Failed to send to external backend".to_string(),
event_clone,
);
self.error_queue.add_error(error_event);
return; }
}
self.events.write(event);
}
pub fn write_batch(&mut self, topic: &str, events: impl IntoIterator<Item = T>) {
let events: Vec<_> = events.into_iter().collect();
if let Some(backend_res) = &self.backend {
let backend = backend_res.read();
for event in &events {
if !backend.try_send(event, topic) {
let error_event = EventBusError::immediate(
topic.to_string(),
crate::EventBusErrorType::Other,
"Failed to send to external backend".to_string(),
event.clone(),
);
self.error_queue.add_error(error_event);
}
}
}
for event in events {
self.events.write(event);
}
}
pub fn write_with_headers(
&mut self,
topic: &str,
event: T,
headers: HashMap<String, String>
) {
let event_clone = event.clone();
if let Some(backend_res) = &self.backend {
let backend = backend_res.read();
if !backend.try_send_with_headers(&event, topic, &headers) {
let error_event = EventBusError::immediate(
topic.to_string(),
crate::EventBusErrorType::Other,
"Failed to send to external backend with headers".to_string(),
event_clone,
);
self.error_queue.add_error(error_event);
return; }
}
self.events.write(event);
}
pub fn write_default(&mut self, topic: &str)
where
T: Default,
{
self.write(topic, T::default())
}
}