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