Skip to main content

sz_orm_queue/
kafka.rs

1use crate::error::MqError;
2use crate::queue::{InMemoryQueue, Message, MessageQueue};
3use async_trait::async_trait;
4
5pub struct InMemoryKafkaQueue {
6    inner: InMemoryQueue,
7}
8
9impl InMemoryKafkaQueue {
10    pub fn new() -> Self {
11        Self {
12            inner: InMemoryQueue::new(),
13        }
14    }
15
16    pub async fn message_count(&self, topic: &str) -> usize {
17        self.inner.message_count(topic).await
18    }
19
20    pub async fn subscriber_count(&self, topic: &str) -> usize {
21        self.inner.subscriber_count(topic).await
22    }
23
24    pub async fn in_flight_count(&self) -> usize {
25        self.inner.in_flight_count().await
26    }
27}
28
29impl Default for InMemoryKafkaQueue {
30    fn default() -> Self {
31        Self::new()
32    }
33}
34
35#[async_trait]
36impl MessageQueue for InMemoryKafkaQueue {
37    async fn publish(&self, topic: &str, message: &[u8]) -> Result<(), MqError> {
38        self.inner.publish(topic, message).await
39    }
40
41    async fn consume(&self, topic: &str) -> Result<Option<Message>, MqError> {
42        self.inner.consume(topic).await
43    }
44
45    async fn ack(&self, message_id: &str) -> Result<(), MqError> {
46        self.inner.ack(message_id).await
47    }
48
49    async fn subscribe(&self, topic: &str) -> Result<(), MqError> {
50        self.inner.subscribe(topic).await
51    }
52}
53
54#[cfg(test)]
55mod tests {
56    use super::*;
57
58    #[tokio::test]
59    async fn test_kafka_publish_and_consume() {
60        let queue = InMemoryKafkaQueue::new();
61        queue.publish("topic-a", b"hello kafka").await.unwrap();
62        assert_eq!(queue.message_count("topic-a").await, 1);
63
64        let msg = queue
65            .consume("topic-a")
66            .await
67            .unwrap()
68            .expect("message should exist");
69        assert_eq!(msg.payload, b"hello kafka");
70        assert_eq!(msg.topic, "topic-a");
71        assert!(!msg.id.is_empty());
72
73        queue.ack(&msg.id).await.unwrap();
74        assert_eq!(queue.in_flight_count().await, 0);
75    }
76
77    #[tokio::test]
78    async fn test_kafka_consume_empty() {
79        let queue = InMemoryKafkaQueue::new();
80        let result = queue.consume("empty-topic").await.unwrap();
81        assert!(result.is_none());
82    }
83
84    #[tokio::test]
85    async fn test_kafka_ack_unknown() {
86        let queue = InMemoryKafkaQueue::new();
87        let result = queue.ack("nonexistent").await;
88        assert!(result.is_err());
89    }
90
91    #[tokio::test]
92    async fn test_kafka_subscribe_increments() {
93        let queue = InMemoryKafkaQueue::new();
94        queue.subscribe("topic-a").await.unwrap();
95        queue.subscribe("topic-a").await.unwrap();
96        queue.subscribe("topic-b").await.unwrap();
97        assert_eq!(queue.subscriber_count("topic-a").await, 2);
98        assert_eq!(queue.subscriber_count("topic-b").await, 1);
99    }
100
101    #[tokio::test]
102    async fn test_kafka_multiple_messages_order() {
103        let queue = InMemoryKafkaQueue::new();
104        queue.publish("topic-a", b"first").await.unwrap();
105        queue.publish("topic-a", b"second").await.unwrap();
106        queue.publish("topic-a", b"third").await.unwrap();
107
108        let m1 = queue.consume("topic-a").await.unwrap().unwrap();
109        let m2 = queue.consume("topic-a").await.unwrap().unwrap();
110        let m3 = queue.consume("topic-a").await.unwrap().unwrap();
111        assert_eq!(m1.payload, b"first");
112        assert_eq!(m2.payload, b"second");
113        assert_eq!(m3.payload, b"third");
114    }
115}