1use crate::error::MqError;
2use crate::queue::{InMemoryQueue, Message, MessageQueue};
3use async_trait::async_trait;
4
5pub struct InMemoryActivemqQueue {
6 inner: InMemoryQueue,
7}
8
9impl InMemoryActivemqQueue {
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 InMemoryActivemqQueue {
30 fn default() -> Self {
31 Self::new()
32 }
33}
34
35#[async_trait]
36impl MessageQueue for InMemoryActivemqQueue {
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_activemq_publish_and_consume() {
60 let queue = InMemoryActivemqQueue::new();
61 queue
62 .publish("active-topic", b"hello active")
63 .await
64 .unwrap();
65 let msg = queue
66 .consume("active-topic")
67 .await
68 .unwrap()
69 .expect("message should exist");
70 assert_eq!(msg.payload, b"hello active");
71 queue.ack(&msg.id).await.unwrap();
72 }
73
74 #[tokio::test]
75 async fn test_activemq_consume_empty() {
76 let queue = InMemoryActivemqQueue::new();
77 assert!(queue.consume("empty").await.unwrap().is_none());
78 }
79
80 #[tokio::test]
81 async fn test_activemq_subscribe() {
82 let queue = InMemoryActivemqQueue::new();
83 queue.subscribe("active-topic").await.unwrap();
84 assert_eq!(queue.subscriber_count("active-topic").await, 1);
85 }
86}