Skip to main content

sz_orm_queue/
lib.rs

1//! # SZ-ORM Queue — 消息队列
2//!
3//! 提供统一的消息队列抽象,内置 InMemory 实现,并支持 RabbitMQ、Kafka、NATS、
4//! ActiveMQ、Pulsar、RocketMQ 等多种消息中间件 provider。
5//!
6//! ## 主要模块
7//!
8//! - [`queue`] — 核心 trait 与统一封装
9//! - [`rabbitmq`] / [`kafka`] / [`nats`] / [`activemq`] / [`pulsar`] / [`rocketmq`] — 各 provider 实现
10
11pub 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// ============================================================================
60// 真实实现(通过 feature flag 启用)
61// ============================================================================
62
63// RabbitMQ: lapin (AMQP 0.9.1) — 真实实现
64#[cfg(feature = "rabbitmq")]
65pub mod lapin_rabbitmq;
66
67#[cfg(feature = "rabbitmq")]
68pub use lapin_rabbitmq::LapinRabbitmqQueue;
69
70// ActiveMQ: lapin (AMQP 1.0,ActiveMQ Artemis) — 真实实现
71#[cfg(feature = "activemq")]
72pub mod real_activemq;
73
74#[cfg(feature = "activemq")]
75pub use real_activemq::RealActivemqQueue;
76
77// NATS: async-nats — 真实实现
78#[cfg(feature = "nats")]
79pub mod real_nats;
80
81#[cfg(feature = "nats")]
82pub use real_nats::RealNatsQueue;
83
84// Pulsar: pulsar crate — 真实实现
85#[cfg(feature = "pulsar")]
86pub mod real_pulsar;
87
88#[cfg(feature = "pulsar")]
89pub use real_pulsar::RealPulsarQueue;
90
91// Kafka: rdkafka — 真实实现
92#[cfg(feature = "kafka")]
93pub mod real_kafka;
94
95#[cfg(feature = "kafka")]
96pub use real_kafka::RealKafkaQueue;
97
98// RocketMQ: 无成熟 Rust 客户端,保持 stub
99// 跟踪项:https://github.com/mxsm/rocketmq-rust (未来可能可用)
100
101#[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        // 验证消息真正进入队列:consume 应能取回,且内容与发布一致
160        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    // ========================================================================
364    // H-3 修复测试:InMemoryQueue 消息大小限制
365    // ========================================================================
366
367    #[tokio::test]
368    async fn test_h3_in_memory_queue_max_messages_limit() {
369        // 设置每 topic 最大 5 条消息
370        let queue = InMemoryQueue::with_max_messages_per_topic(5);
371
372        // 前 5 条成功
373        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        // 第 6 条应失败
382        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        // 消费后可继续发布
389        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        // 默认 100,000 限制,正常使用不应触发
397        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        // 不同 topic 独立计数
410        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        // topic-a 满了
414        assert!(queue.publish("topic-a", b"a3").await.is_err());
415        // topic-b 不受影响
416        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}