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,
}
enum InstalledApp {
#[allow(dead_code)]
Orders,
}
#[producer(topic = "orders", name = "create_order")]
pub async fn create_order(id: u64) -> Result<Order, StreamingError> {
Ok(Order { id })
}
#[consumer(topic = "orders", group = "processor", name = "handle_order")]
pub async fn handle_order(_msg: Message<Order>) -> Result<(), StreamingError> {
Ok(())
}
#[streaming_patterns(InstalledApp::Orders)]
pub fn app_streaming_routes() -> reinhardt_streaming::StreamingRouter {
streaming_routes![create_order, handle_order]
}
#[test]
fn streaming_patterns_struct_has_correct_topics() {
let urls = OrdersStreamingUrls {
_marker: ::core::marker::PhantomData,
};
assert_eq!(urls.create_order(), "orders");
assert_eq!(urls.handle_order(), "orders");
}