1pub mod error;
12pub mod queue;
13
14pub mod activemq;
15pub mod kafka;
16pub mod nats;
17pub mod pulsar;
18pub mod rabbitmq;
19pub mod rocketmq;
20
21pub use error::MqError;
22
23pub use queue::ActiveConfig;
24pub use queue::InMemoryQueue;
25pub use queue::KafkaConfig;
26pub use queue::Message;
27pub use queue::MessageQueue;
28pub use queue::MqProvider;
29pub use queue::NatsConfig;
30pub use queue::PulsarConfig;
31pub use queue::QueueConfig;
32pub use queue::QueueWrapper;
33pub use queue::RabbitConfig;
34pub use queue::RocketConfig;
35
36pub use activemq::InMemoryActivemqQueue;
37pub use kafka::InMemoryKafkaQueue;
38pub use nats::InMemoryNatsQueue;
39pub use pulsar::InMemoryPulsarQueue;
40pub use rabbitmq::InMemoryRabbitmqQueue;
41pub use rocketmq::InMemoryRocketmqQueue;
42
43#[cfg(feature = "rabbitmq")]
49pub mod lapin_rabbitmq;
50
51#[cfg(feature = "rabbitmq")]
52pub use lapin_rabbitmq::LapinRabbitmqQueue;
53
54#[cfg(feature = "activemq")]
56pub mod real_activemq;
57
58#[cfg(feature = "activemq")]
59pub use real_activemq::RealActivemqQueue;
60
61#[cfg(feature = "nats")]
63pub mod real_nats;
64
65#[cfg(feature = "nats")]
66pub use real_nats::RealNatsQueue;
67
68#[cfg(feature = "pulsar")]
70pub mod real_pulsar;
71
72#[cfg(feature = "pulsar")]
73pub use real_pulsar::RealPulsarQueue;
74
75#[cfg(feature = "kafka")]
77pub mod real_kafka;
78
79#[cfg(feature = "kafka")]
80pub use real_kafka::RealKafkaQueue;
81
82#[cfg(test)]
86mod tests {
87 use super::*;
88
89 #[test]
90 fn test_message_new() {
91 let msg = Message::new("test", vec![1, 2, 3]);
92 assert_eq!(msg.topic, "test");
93 assert_eq!(msg.payload, vec![1, 2, 3]);
94 assert!(msg.key.is_none());
95 }
96
97 #[test]
98 fn test_message_with_key() {
99 let msg = Message::new("test", vec![]).with_key("mykey");
100 assert_eq!(msg.key, Some("mykey".to_string()));
101 }
102
103 #[test]
104 fn test_message_text() {
105 let msg = Message::text_message("test", "hello");
106 assert_eq!(msg.text(), Some("hello"));
107 }
108
109 #[test]
110 fn test_message_json() {
111 let msg = Message::json_message("test", &serde_json::json!({"key": "value"})).unwrap();
112 let parsed: serde_json::Value = msg.json().unwrap();
113 assert_eq!(parsed["key"], "value");
114 }
115
116 #[test]
117 fn test_queue_config_default() {
118 let config = QueueConfig::default();
119 assert_eq!(config.brokers, vec!["localhost:9092".to_string()]);
120 assert!(matches!(config.provider, MqProvider::Kafka(_)));
121 }
122
123 #[test]
124 fn test_queue_config_builder() {
125 let config = QueueConfig::new()
126 .with_provider(MqProvider::RabbitMQ(RabbitConfig::default()))
127 .with_brokers(vec!["localhost:5672".to_string()])
128 .with_group("my-group")
129 .with_auth("user", "pass");
130
131 assert!(matches!(config.provider, MqProvider::RabbitMQ(_)));
132 assert_eq!(config.brokers, vec!["localhost:5672".to_string()]);
133 assert_eq!(config.group_id, Some("my-group".to_string()));
134 assert_eq!(config.username, Some("user".to_string()));
135 assert_eq!(config.password, Some("pass".to_string()));
136 }
137
138 #[tokio::test]
139 async fn test_queue_wrapper_publish() {
140 let wrapper = QueueWrapper::new(MqProvider::Kafka(KafkaConfig::default()));
141 let result = wrapper.publish("test-topic", b"message").await;
142 assert!(result.is_ok());
143 }
144
145 #[test]
146 fn test_kafka_config_default() {
147 let config = KafkaConfig::default();
148 assert!(config.client_id.is_none());
149 assert!(config.acks.is_none());
150 }
151
152 #[test]
153 fn test_rabbit_config_default() {
154 let config = RabbitConfig::default();
155 assert!(config.virtual_host.is_none());
156 }
157
158 #[test]
159 fn test_rocket_config_default() {
160 let config = RocketConfig::default();
161 assert!(config.namespace.is_none());
162 }
163
164 #[test]
165 fn test_active_config_default() {
166 let config = ActiveConfig::default();
167 assert!(config.broker_url.is_none());
168 }
169
170 #[test]
171 fn test_nats_config_default() {
172 let config = NatsConfig::default();
173 assert!(config.name.is_none());
174 }
175
176 #[test]
177 fn test_pulsar_config_default() {
178 let config = PulsarConfig::default();
179 assert!(config.service_url.is_none());
180 }
181
182 #[test]
183 fn test_message_timestamp_set() {
184 let msg = Message::new("test", vec![]);
185 assert!(msg.timestamp > 0);
186 }
187
188 #[test]
189 fn test_message_headers() {
190 let msg = Message::new("test", vec![]);
191 assert!(msg.headers.is_empty());
192 }
193
194 #[tokio::test]
195 async fn test_in_memory_queue_publish_and_consume() {
196 let queue = InMemoryQueue::new();
197 queue.publish("orders", b"order-1").await.unwrap();
198 assert_eq!(queue.message_count("orders").await, 1);
199
200 let msg = queue
201 .consume("orders")
202 .await
203 .unwrap()
204 .expect("message should exist");
205 assert_eq!(msg.payload, b"order-1");
206 assert_eq!(msg.topic, "orders");
207 assert!(!msg.id.is_empty());
208 assert_eq!(queue.message_count("orders").await, 0);
209 }
210
211 #[tokio::test]
212 async fn test_in_memory_queue_consume_empty() {
213 let queue = InMemoryQueue::new();
214 let result = queue.consume("empty-topic").await.unwrap();
215 assert!(result.is_none());
216 }
217
218 #[tokio::test]
219 async fn test_in_memory_queue_ack() {
220 let queue = InMemoryQueue::new();
221 queue.publish("topic", b"data").await.unwrap();
222 let msg = queue.consume("topic").await.unwrap().unwrap();
223
224 assert_eq!(queue.in_flight_count().await, 1);
225 queue.ack(&msg.id).await.unwrap();
226 assert_eq!(queue.in_flight_count().await, 0);
227 }
228
229 #[tokio::test]
230 async fn test_in_memory_queue_ack_unknown_id() {
231 let queue = InMemoryQueue::new();
232 let result = queue.ack("unknown-id").await;
233 assert!(result.is_err());
234 }
235
236 #[tokio::test]
237 async fn test_in_memory_queue_subscribe() {
238 let queue = InMemoryQueue::new();
239 queue.subscribe("topic-a").await.unwrap();
240 queue.subscribe("topic-a").await.unwrap();
241 queue.subscribe("topic-b").await.unwrap();
242 assert_eq!(queue.subscriber_count("topic-a").await, 2);
243 assert_eq!(queue.subscriber_count("topic-b").await, 1);
244 assert_eq!(queue.subscriber_count("topic-c").await, 0);
245 }
246
247 #[tokio::test]
248 async fn test_in_memory_queue_multiple_topics() {
249 let queue = InMemoryQueue::new();
250 queue.publish("topic-a", b"a1").await.unwrap();
251 queue.publish("topic-b", b"b1").await.unwrap();
252 queue.publish("topic-a", b"a2").await.unwrap();
253
254 assert_eq!(queue.message_count("topic-a").await, 2);
255 assert_eq!(queue.message_count("topic-b").await, 1);
256 }
257
258 #[tokio::test]
259 async fn test_in_memory_queue_fifo_order() {
260 let queue = InMemoryQueue::new();
261 queue.publish("topic", b"first").await.unwrap();
262 queue.publish("topic", b"second").await.unwrap();
263 queue.publish("topic", b"third").await.unwrap();
264
265 let m1 = queue.consume("topic").await.unwrap().unwrap();
266 let m2 = queue.consume("topic").await.unwrap().unwrap();
267 let m3 = queue.consume("topic").await.unwrap().unwrap();
268 assert_eq!(m1.payload, b"first");
269 assert_eq!(m2.payload, b"second");
270 assert_eq!(m3.payload, b"third");
271 }
272
273 #[tokio::test]
274 async fn test_queue_wrapper_publish_and_consume() {
275 let wrapper = QueueWrapper::new(MqProvider::Kafka(KafkaConfig::default()));
276 wrapper.publish("wrapper-topic", b"payload").await.unwrap();
277
278 let msg = wrapper
279 .consume("wrapper-topic")
280 .await
281 .unwrap()
282 .expect("message should exist");
283 assert_eq!(msg.payload, b"payload");
284 wrapper.ack(&msg.id).await.unwrap();
285 }
286
287 #[tokio::test]
288 async fn test_queue_wrapper_subscribe_and_ack() {
289 let wrapper = QueueWrapper::new(MqProvider::RabbitMQ(RabbitConfig::default()));
290 wrapper.subscribe("sub-topic").await.unwrap();
291 wrapper.publish("sub-topic", b"hello").await.unwrap();
292
293 let msg = wrapper
294 .consume("sub-topic")
295 .await
296 .unwrap()
297 .expect("message should exist");
298 assert_eq!(msg.payload, b"hello");
299 assert!(wrapper.ack(&msg.id).await.is_ok());
300 }
301
302 #[tokio::test]
303 async fn test_queue_wrapper_consume_empty() {
304 let wrapper = QueueWrapper::new(MqProvider::Nats(NatsConfig::default()));
305 let result = wrapper.consume("empty").await.unwrap();
306 assert!(result.is_none());
307 }
308
309 #[tokio::test]
310 async fn test_queue_wrapper_ack_unknown() {
311 let wrapper = QueueWrapper::new(MqProvider::Pulsar(PulsarConfig::default()));
312 assert!(wrapper.ack("unknown").await.is_err());
313 }
314
315 #[tokio::test]
316 async fn test_queue_wrapper_with_each_provider() {
317 let providers = vec![
318 MqProvider::Kafka(KafkaConfig::default()),
319 MqProvider::RabbitMQ(RabbitConfig::default()),
320 MqProvider::RocketMQ(RocketConfig::default()),
321 MqProvider::ActiveMQ(ActiveConfig::default()),
322 MqProvider::Nats(NatsConfig::default()),
323 MqProvider::Pulsar(PulsarConfig::default()),
324 ];
325
326 for provider in providers {
327 let wrapper = QueueWrapper::new(provider);
328 wrapper.publish("topic", b"data").await.unwrap();
329 let msg = wrapper
330 .consume("topic")
331 .await
332 .unwrap()
333 .expect("message should exist");
334 assert_eq!(msg.payload, b"data");
335 wrapper.ack(&msg.id).await.unwrap();
336 }
337 }
338
339 #[tokio::test]
344 async fn test_h3_in_memory_queue_max_messages_limit() {
345 let queue = InMemoryQueue::with_max_messages_per_topic(5);
347
348 for i in 0..5 {
350 queue
351 .publish("limited", format!("msg-{i}").as_bytes())
352 .await
353 .unwrap();
354 }
355 assert_eq!(queue.message_count("limited").await, 5);
356
357 let result = queue.publish("limited", b"overflow").await;
359 assert!(result.is_err());
360 let err_msg = result.unwrap_err().to_string();
361 assert!(err_msg.contains("H-3 protection"), "err: {err_msg}");
362 assert!(err_msg.contains("limited"), "err: {err_msg}");
363
364 let _ = queue.consume("limited").await.unwrap().unwrap();
366 queue.publish("limited", b"after-consume").await.unwrap();
367 assert_eq!(queue.message_count("limited").await, 5);
368 }
369
370 #[tokio::test]
371 async fn test_h3_in_memory_queue_default_limit_accepts_normal_usage() {
372 let queue = InMemoryQueue::new();
374 for i in 0..1000 {
375 queue
376 .publish("normal", format!("msg-{i}").as_bytes())
377 .await
378 .unwrap();
379 }
380 assert_eq!(queue.message_count("normal").await, 1000);
381 }
382
383 #[tokio::test]
384 async fn test_h3_in_memory_queue_limit_isolated_per_topic() {
385 let queue = InMemoryQueue::with_max_messages_per_topic(2);
387 queue.publish("topic-a", b"a1").await.unwrap();
388 queue.publish("topic-a", b"a2").await.unwrap();
389 assert!(queue.publish("topic-a", b"a3").await.is_err());
391 queue.publish("topic-b", b"b1").await.unwrap();
393 queue.publish("topic-b", b"b2").await.unwrap();
394 assert!(queue.publish("topic-b", b"b3").await.is_err());
395 }
396}