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