use bevy::ecs::system::SystemParam;
use bevy::prelude::*;
use std::marker::PhantomData;
use bevy_event_bus::backends::event_bus_backend::{LagReportingDescriptor, ManualCommitDescriptor};
use bevy_event_bus::backends::{EventBusBackend, EventBusBackendResource};
use bevy_event_bus::decoder::DecoderRegistry;
use bevy_event_bus::resources::{
ConsumerMetrics, DecodedEventBuffer, DrainMetricsEvent, DrainedTopicMetadata,
EventBusConsumerConfig, EventMetadata, IncomingMessage, MessageQueue, ProcessedMessage,
ProvisionedTopology,
};
use bevy_event_bus::runtime::{block_on, ensure_runtime};
use bevy_event_bus::writers::EventBusErrorQueue;
pub struct EventBusPlugin;
impl Plugin for EventBusPlugin {
fn build(&self, app: &mut App) {
app.add_event::<bevy_event_bus::EventBusDecodeError>();
}
}
pub struct EventBusPlugins<B: EventBusBackend>(pub B);
#[derive(Event, Debug, Clone)]
pub struct BackendReadyEvent {
pub backend: String,
pub topics: Vec<String>,
}
#[derive(Event, Debug, Clone)]
pub struct BackendDownEvent {
pub backend: String,
pub reason: String,
}
#[derive(Debug)]
enum LifecycleMessage {
Ready {
backend: String,
topics: Vec<String>,
},
}
#[derive(Resource)]
struct BackendLifecycleChannel(crossbeam_channel::Receiver<LifecycleMessage>);
#[derive(Resource, Debug, Clone, Copy)]
pub struct BackendStatus {
pub ready: bool,
}
#[derive(Resource, Debug, Clone)]
pub struct BackendCapabilities {
pub backend: String,
pub message_stream: bool,
pub manual_commit: Option<ManualCommitDescriptor>,
pub lag_reporting: Option<LagReportingDescriptor>,
}
impl BackendCapabilities {
fn new(backend: impl Into<String>) -> Self {
Self {
backend: backend.into(),
message_stream: false,
manual_commit: None,
lag_reporting: None,
}
}
}
#[derive(SystemParam)]
struct DrainSystemResources<'w, 's> {
metadata_buffers: ResMut<'w, DrainedTopicMetadata>,
decoded_buffer: ResMut<'w, DecodedEventBuffer>,
decoder_registry: ResMut<'w, DecoderRegistry>,
metrics: ResMut<'w, ConsumerMetrics>,
config: Res<'w, EventBusConsumerConfig>,
_marker: PhantomData<&'s ()>,
}
impl<B: EventBusBackend> Plugin for EventBusPlugins<B> {
fn build(&self, app: &mut App) {
app.add_plugins(EventBusPlugin);
ensure_runtime(app);
self.0.configure_plugin(app);
let boxed = self.0.clone_box();
app.insert_resource(EventBusBackendResource::from_box(boxed));
app.insert_resource(EventBusErrorQueue::default());
app.init_resource::<ProvisionedTopology>();
self.0.apply_event_bindings(app);
bevy_event_bus::writers::outbound_bridge::activate_registered_bridges(app);
if let Some(backend_res) = app
.world()
.get_resource::<EventBusBackendResource>()
.cloned()
{
let (backend_name, mut backend_setup) = {
let mut guard = backend_res.write();
let _ = block_on(guard.connect());
let backend_name = guard.backend_name().to_string();
let setup = guard.setup_plugin(app.world_mut());
(backend_name, setup)
};
let mut capabilities = BackendCapabilities::new(backend_name.clone());
if let Some(topology) = backend_setup.kafka_topology.take() {
let mut registry = app.world_mut().resource_mut::<ProvisionedTopology>();
registry.record_kafka(topology);
}
if let Some(topology) = backend_setup.redis_topology.take() {
let mut registry = app.world_mut().resource_mut::<ProvisionedTopology>();
registry.record_redis(topology);
}
if let Some(rx) = backend_setup.message_stream.take() {
app.world_mut()
.insert_resource(MessageQueue { receiver: rx });
capabilities.message_stream = true;
}
if let Some(handle) = backend_setup.manual_commit.take() {
let descriptor = handle.descriptor();
{
let world = app.world_mut();
handle.register_resources(world);
}
capabilities.manual_commit = Some(descriptor);
}
if let Some(handle) = backend_setup.lag_reporting.take() {
let descriptor = handle.descriptor();
{
let world = app.world_mut();
handle.register_resources(world);
}
capabilities.lag_reporting = Some(descriptor);
}
let (tx, rx_life) = crossbeam_channel::unbounded::<LifecycleMessage>();
let ready_topics = backend_setup.ready_topics.clone();
let _ = tx.send(LifecycleMessage::Ready {
backend: backend_name.clone(),
topics: ready_topics.clone(),
});
app.world_mut()
.insert_resource(BackendLifecycleChannel(rx_life));
app.world_mut()
.insert_resource(BackendStatus { ready: false });
app.world_mut().insert_resource(capabilities);
}
app.init_resource::<DrainedTopicMetadata>();
app.init_resource::<DecodedEventBuffer>();
app.init_resource::<DecoderRegistry>();
app.init_resource::<ConsumerMetrics>();
app.init_resource::<EventBusConsumerConfig>();
app.add_event::<DrainMetricsEvent>();
app.add_event::<BackendReadyEvent>();
fn drain_system(
backend: Option<Res<EventBusBackendResource>>,
resources: DrainSystemResources,
maybe_queue: Option<Res<MessageQueue>>,
mut drain_events: EventWriter<DrainMetricsEvent>,
) {
let DrainSystemResources {
metadata_buffers,
decoded_buffer,
decoder_registry,
metrics,
config,
..
} = resources;
let mut metadata_buffers = metadata_buffers;
let mut decoded_buffer = decoded_buffer;
let mut decoder_registry = decoder_registry;
let mut metrics = metrics;
let frame_start = std::time::Instant::now();
metrics.drained_last_frame = 0;
metrics.queue_len_start = 0;
metrics.queue_len_end = 0;
let backend_ref = backend.as_ref();
if let Some(queue) = maybe_queue {
metrics.queue_len_start = queue.receiver.len();
let limit = config.max_events_per_frame;
let time_budget = config
.max_drain_millis
.map(std::time::Duration::from_millis);
while limit
.map(|l| metrics.drained_last_frame < l)
.unwrap_or(true)
{
if let Some(budget) = time_budget {
if frame_start.elapsed() >= budget {
break;
}
if budget <= std::time::Duration::from_millis(1)
&& metrics.drained_last_frame > 0
{
break;
}
}
match queue.receiver.try_recv() {
Ok(msg) => {
let IncomingMessage {
source,
payload,
key,
timestamp,
backend_metadata,
} = msg;
let topic = source;
let topic_str = topic.as_str();
tracing::debug!(topic=%topic_str, "Processing message with multi-decoder pipeline");
let metadata =
EventMetadata::new(topic.clone(), timestamp, key, backend_metadata);
let decoded_events = decoder_registry.decode_all(topic_str, &payload);
let topic_buffer =
decoded_buffer.topics.entry(topic.clone()).or_default();
topic_buffer.total_processed += 1;
if decoded_events.is_empty() {
topic_buffer.decode_failures += 1;
let decoder_count = decoder_registry.decoder_count(topic_str);
tracing::debug!(
topic = %topic_str,
decoders_tried = decoder_count,
"No decoder succeeded for message"
);
let decode_error = bevy_event_bus::EventBusDecodeError::new(
metadata.source.clone(),
format!(
"No decoder succeeded. Tried {} decoders",
decoder_count
),
payload,
format!("tried_{}_decoders", decoder_count),
Some(metadata.clone()),
);
metadata_buffers.decode_errors.push(decode_error);
} else {
for decoded_event in decoded_events {
tracing::trace!(
topic = %topic_str,
decoder = %decoded_event.decoder_name,
"Successfully decoded event"
);
let type_erased = crate::resources::TypeErasedEvent {
event: decoded_event.event,
metadata: metadata.clone(),
decoder_name: decoded_event.decoder_name,
};
topic_buffer
.events_by_type
.entry(decoded_event.type_id)
.or_default()
.push(type_erased);
}
let processed_msg = ProcessedMessage {
payload,
metadata: metadata.clone(),
};
metadata_buffers
.topics
.entry(topic.clone())
.or_default()
.push(processed_msg);
}
metrics.drained_last_frame += 1;
}
Err(crossbeam_channel::TryRecvError::Empty) => break,
Err(crossbeam_channel::TryRecvError::Disconnected) => break,
}
}
metrics.remaining_channel_after_drain = queue.receiver.len();
metrics.queue_len_end = metrics.remaining_channel_after_drain;
metrics.total_drained += metrics.drained_last_frame;
}
if let Some(backend_res) = backend_ref {
let guard = backend_res.read();
guard.augment_metrics(metrics.as_mut());
}
if metrics.drained_last_frame == 0 {
metrics.idle_frames += 1;
}
if metrics.idle_frames % 30 == 0 && metrics.idle_frames > 0 {
let before_count = metadata_buffers.topics.len();
metadata_buffers
.topics
.retain(|_topic, buffer| !buffer.is_empty());
let after_count = metadata_buffers.topics.len();
if before_count > after_count {
tracing::debug!(
cleaned_topics = before_count - after_count,
remaining_topics = after_count,
"Cleaned up empty topic metadata buffers"
);
}
let before_decoded = decoded_buffer.topics.len();
decoded_buffer
.topics
.retain(|_topic, buffer| buffer.total_events() > 0);
let after_decoded = decoded_buffer.topics.len();
if before_decoded > after_decoded {
tracing::debug!(
cleaned_decoded_topics = before_decoded - after_decoded,
remaining_decoded_topics = after_decoded,
"Cleaned up empty decoded event buffers"
);
}
}
metrics.drain_duration_us = frame_start.elapsed().as_micros();
if metrics.drain_duration_us == 0
&& (metrics.drained_last_frame > 0 || metrics.queue_len_start > 0)
{
metrics.drain_duration_us = 1; }
drain_events.write(DrainMetricsEvent {
drained: metrics.drained_last_frame,
remaining: metrics.remaining_channel_after_drain,
total_drained: metrics.total_drained,
dropped: metrics.dropped_messages,
drain_duration_us: metrics.drain_duration_us,
});
}
app.add_systems(PreUpdate, drain_system);
app.add_systems(PreUpdate, decode_error_dispatch_system.after(drain_system));
fn error_queue_flush_system(world: &mut World) {
let pending_errors = {
let error_queue = world.resource::<EventBusErrorQueue>();
error_queue.drain_pending()
};
for error_fn in pending_errors {
error_fn(world);
}
}
app.add_systems(PostUpdate, error_queue_flush_system);
fn lifecycle_system(
maybe_channel: Option<Res<BackendLifecycleChannel>>,
mut ready_writer: EventWriter<BackendReadyEvent>,
) {
if let Some(ch) = maybe_channel {
while let Ok(msg) = ch.0.try_recv() {
match msg {
LifecycleMessage::Ready { backend, topics } => {
ready_writer.write(BackendReadyEvent {
backend: backend.clone(),
topics: topics.clone(),
});
}
}
}
}
}
app.add_systems(PreUpdate, lifecycle_system);
fn backend_status_update(
status: Option<ResMut<BackendStatus>>,
mut ready_events: EventReader<BackendReadyEvent>,
) {
if let Some(mut s) = status {
for _ev in ready_events.read() {
s.ready = true;
}
}
}
app.add_systems(PreUpdate, backend_status_update);
fn decode_error_dispatch_system(
mut drained_metadata: ResMut<DrainedTopicMetadata>,
mut decode_error_writer: EventWriter<bevy_event_bus::EventBusDecodeError>,
) {
for decode_error in drained_metadata.decode_errors.drain(..) {
decode_error_writer.write(decode_error);
}
}
}
}