pub mod backend;
pub mod error;
pub mod in_memory;
pub mod message;
#[cfg(feature = "kafka")]
pub mod kafka;
#[cfg(feature = "kafka")]
pub mod di;
#[cfg(feature = "kafka")]
pub mod global;
#[cfg(feature = "kafka")]
pub use global::{global_producer, set_global_producer};
pub mod macros;
pub mod router;
pub use router::{
ConsumerFactory, StreamingHandlerKind, StreamingHandlerRegistration, StreamingRouter,
};
pub trait StreamingTopicResolver {
fn resolve_topic(&self, name: &str) -> &'static str;
}
pub use backend::StreamingBackend;
pub use error::StreamingError;
pub use in_memory::InMemoryStreamingBackend;
pub use message::Message;
#[derive(Debug)]
pub struct StreamingHandlerMetadata {
pub name: &'static str,
pub topic: &'static str,
pub kind: StreamingHandlerKind,
pub group: Option<&'static str>,
pub module_path: &'static str,
}
inventory::collect!(StreamingHandlerMetadata);
pub fn resolve_streaming_topic(name: &str) -> &'static str {
for meta in inventory::iter::<StreamingHandlerMetadata> {
if meta.name == name {
return meta.topic;
}
}
panic!(
"Streaming handler `{name}` not registered. Ensure the function is annotated with `#[producer]` or `#[consumer]`."
);
}