use core::marker::PhantomData;
use reinhardt_macros::{consumer, producer, streaming_patterns};
use reinhardt_streaming::{Message, StreamingError, streaming_routes};
use serde::{Deserialize, Serialize};
#[derive(Debug, Serialize, Deserialize, PartialEq)]
struct Order {
id: u64,
}
struct AppLabel;
#[producer(topic = "events", name = "emit_event")]
pub async fn emit_event(id: u64) -> Result<Order, StreamingError> {
Ok(Order { id })
}
#[consumer(topic = "events", group = "listener", name = "on_event")]
pub async fn on_event(_msg: Message<Order>) -> Result<(), StreamingError> {
Ok(())
}
#[streaming_patterns(AppLabel)]
pub fn streaming_routes_fn() -> reinhardt_streaming::StreamingRouter {
streaming_routes![emit_event, on_event]
}
#[test]
fn per_app_struct_has_correct_topic_names() {
let app_urls = ApplabelStreamingUrls {
_marker: PhantomData,
};
assert_eq!(app_urls.emit_event(), "events");
assert_eq!(app_urls.on_event(), "events");
}
#[test]
fn streaming_resolvers_module_generated() {
#[allow(unused_imports)]
use streaming_resolvers::__streaming_resolver_meta_emit_event;
}