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 let msg = wrapper
149 .consume("test-topic")
150 .await
151 .expect("consume 不应报错")
152 .expect("消息应存在");
153 assert_eq!(msg.payload, b"message", "取回的 payload 应与发布一致");
154 assert_eq!(msg.topic, "test-topic", "取回的 topic 应与发布一致");
155 }
156
157 #[test]
158 fn test_kafka_config_default() {
159 let config = KafkaConfig::default();
160 assert!(config.client_id.is_none());
161 assert!(config.acks.is_none());
162 }
163
164 #[test]
165 fn test_rabbit_config_default() {
166 let config = RabbitConfig::default();
167 assert!(config.virtual_host.is_none());
168 }
169
170 #[test]
171 fn test_rocket_config_default() {
172 let config = RocketConfig::default();
173 assert!(config.namespace.is_none());
174 }
175
176 #[test]
177 fn test_active_config_default() {
178 let config = ActiveConfig::default();
179 assert!(config.broker_url.is_none());
180 }
181
182 #[test]
183 fn test_nats_config_default() {
184 let config = NatsConfig::default();
185 assert!(config.name.is_none());
186 }
187
188 #[test]
189 fn test_pulsar_config_default() {
190 let config = PulsarConfig::default();
191 assert!(config.service_url.is_none());
192 }
193
194 #[test]
195 fn test_message_timestamp_set() {
196 let msg = Message::new("test", vec![]);
197 assert!(msg.timestamp > 0);
198 }
199
200 #[test]
201 fn test_message_headers() {
202 let msg = Message::new("test", vec![]);
203 assert!(msg.headers.is_empty());
204 }
205
206 #[tokio::test]
207 async fn test_in_memory_queue_publish_and_consume() {
208 let queue = InMemoryQueue::new();
209 queue.publish("orders", b"order-1").await.unwrap();
210 assert_eq!(queue.message_count("orders").await, 1);
211
212 let msg = queue
213 .consume("orders")
214 .await
215 .unwrap()
216 .expect("message should exist");
217 assert_eq!(msg.payload, b"order-1");
218 assert_eq!(msg.topic, "orders");
219 assert!(!msg.id.is_empty());
220 assert_eq!(queue.message_count("orders").await, 0);
221 }
222
223 #[tokio::test]
224 async fn test_in_memory_queue_consume_empty() {
225 let queue = InMemoryQueue::new();
226 let result = queue.consume("empty-topic").await.unwrap();
227 assert!(result.is_none());
228 }
229
230 #[tokio::test]
231 async fn test_in_memory_queue_ack() {
232 let queue = InMemoryQueue::new();
233 queue.publish("topic", b"data").await.unwrap();
234 let msg = queue.consume("topic").await.unwrap().unwrap();
235
236 assert_eq!(queue.in_flight_count().await, 1);
237 queue.ack(&msg.id).await.unwrap();
238 assert_eq!(queue.in_flight_count().await, 0);
239 }
240
241 #[tokio::test]
242 async fn test_in_memory_queue_ack_unknown_id() {
243 let queue = InMemoryQueue::new();
244 let result = queue.ack("unknown-id").await;
245 assert!(result.is_err());
246 }
247
248 #[tokio::test]
249 async fn test_in_memory_queue_subscribe() {
250 let queue = InMemoryQueue::new();
251 queue.subscribe("topic-a").await.unwrap();
252 queue.subscribe("topic-a").await.unwrap();
253 queue.subscribe("topic-b").await.unwrap();
254 assert_eq!(queue.subscriber_count("topic-a").await, 2);
255 assert_eq!(queue.subscriber_count("topic-b").await, 1);
256 assert_eq!(queue.subscriber_count("topic-c").await, 0);
257 }
258
259 #[tokio::test]
260 async fn test_in_memory_queue_multiple_topics() {
261 let queue = InMemoryQueue::new();
262 queue.publish("topic-a", b"a1").await.unwrap();
263 queue.publish("topic-b", b"b1").await.unwrap();
264 queue.publish("topic-a", b"a2").await.unwrap();
265
266 assert_eq!(queue.message_count("topic-a").await, 2);
267 assert_eq!(queue.message_count("topic-b").await, 1);
268 }
269
270 #[tokio::test]
271 async fn test_in_memory_queue_fifo_order() {
272 let queue = InMemoryQueue::new();
273 queue.publish("topic", b"first").await.unwrap();
274 queue.publish("topic", b"second").await.unwrap();
275 queue.publish("topic", b"third").await.unwrap();
276
277 let m1 = queue.consume("topic").await.unwrap().unwrap();
278 let m2 = queue.consume("topic").await.unwrap().unwrap();
279 let m3 = queue.consume("topic").await.unwrap().unwrap();
280 assert_eq!(m1.payload, b"first");
281 assert_eq!(m2.payload, b"second");
282 assert_eq!(m3.payload, b"third");
283 }
284
285 #[tokio::test]
286 async fn test_queue_wrapper_publish_and_consume() {
287 let wrapper = QueueWrapper::new(MqProvider::Kafka(KafkaConfig::default()));
288 wrapper.publish("wrapper-topic", b"payload").await.unwrap();
289
290 let msg = wrapper
291 .consume("wrapper-topic")
292 .await
293 .unwrap()
294 .expect("message should exist");
295 assert_eq!(msg.payload, b"payload");
296 wrapper.ack(&msg.id).await.unwrap();
297 }
298
299 #[tokio::test]
300 async fn test_queue_wrapper_subscribe_and_ack() {
301 let wrapper = QueueWrapper::new(MqProvider::RabbitMQ(RabbitConfig::default()));
302 wrapper.subscribe("sub-topic").await.unwrap();
303 wrapper.publish("sub-topic", b"hello").await.unwrap();
304
305 let msg = wrapper
306 .consume("sub-topic")
307 .await
308 .unwrap()
309 .expect("message should exist");
310 assert_eq!(msg.payload, b"hello");
311 assert!(wrapper.ack(&msg.id).await.is_ok());
312 }
313
314 #[tokio::test]
315 async fn test_queue_wrapper_consume_empty() {
316 let wrapper = QueueWrapper::new(MqProvider::Nats(NatsConfig::default()));
317 let result = wrapper.consume("empty").await.unwrap();
318 assert!(result.is_none());
319 }
320
321 #[tokio::test]
322 async fn test_queue_wrapper_ack_unknown() {
323 let wrapper = QueueWrapper::new(MqProvider::Pulsar(PulsarConfig::default()));
324 assert!(wrapper.ack("unknown").await.is_err());
325 }
326
327 #[tokio::test]
328 async fn test_queue_wrapper_with_each_provider() {
329 let providers = vec![
330 MqProvider::Kafka(KafkaConfig::default()),
331 MqProvider::RabbitMQ(RabbitConfig::default()),
332 MqProvider::RocketMQ(RocketConfig::default()),
333 MqProvider::ActiveMQ(ActiveConfig::default()),
334 MqProvider::Nats(NatsConfig::default()),
335 MqProvider::Pulsar(PulsarConfig::default()),
336 ];
337
338 for provider in providers {
339 let wrapper = QueueWrapper::new(provider);
340 wrapper.publish("topic", b"data").await.unwrap();
341 let msg = wrapper
342 .consume("topic")
343 .await
344 .unwrap()
345 .expect("message should exist");
346 assert_eq!(msg.payload, b"data");
347 wrapper.ack(&msg.id).await.unwrap();
348 }
349 }
350
351 #[tokio::test]
356 async fn test_h3_in_memory_queue_max_messages_limit() {
357 let queue = InMemoryQueue::with_max_messages_per_topic(5);
359
360 for i in 0..5 {
362 queue
363 .publish("limited", format!("msg-{i}").as_bytes())
364 .await
365 .unwrap();
366 }
367 assert_eq!(queue.message_count("limited").await, 5);
368
369 let result = queue.publish("limited", b"overflow").await;
371 assert!(result.is_err());
372 let err_msg = result.unwrap_err().to_string();
373 assert!(err_msg.contains("H-3 protection"), "err: {err_msg}");
374 assert!(err_msg.contains("limited"), "err: {err_msg}");
375
376 let _ = queue.consume("limited").await.unwrap().unwrap();
378 queue.publish("limited", b"after-consume").await.unwrap();
379 assert_eq!(queue.message_count("limited").await, 5);
380 }
381
382 #[tokio::test]
383 async fn test_h3_in_memory_queue_default_limit_accepts_normal_usage() {
384 let queue = InMemoryQueue::new();
386 for i in 0..1000 {
387 queue
388 .publish("normal", format!("msg-{i}").as_bytes())
389 .await
390 .unwrap();
391 }
392 assert_eq!(queue.message_count("normal").await, 1000);
393 }
394
395 #[tokio::test]
396 async fn test_h3_in_memory_queue_limit_isolated_per_topic() {
397 let queue = InMemoryQueue::with_max_messages_per_topic(2);
399 queue.publish("topic-a", b"a1").await.unwrap();
400 queue.publish("topic-a", b"a2").await.unwrap();
401 assert!(queue.publish("topic-a", b"a3").await.is_err());
403 queue.publish("topic-b", b"b1").await.unwrap();
405 queue.publish("topic-b", b"b2").await.unwrap();
406 assert!(queue.publish("topic-b", b"b3").await.is_err());
407 }
408}