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
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// ============================================================================
48// 真实实现(通过 feature flag 启用)
49// ============================================================================
50
51// RabbitMQ: lapin (AMQP 0.9.1) — 真实实现
52#[cfg(feature = "rabbitmq")]
53pub mod lapin_rabbitmq;
54
55#[cfg(feature = "rabbitmq")]
56pub use lapin_rabbitmq::LapinRabbitmqQueue;
57
58// ActiveMQ: lapin (AMQP 1.0,ActiveMQ Artemis) — 真实实现
59#[cfg(feature = "activemq")]
60pub mod real_activemq;
61
62#[cfg(feature = "activemq")]
63pub use real_activemq::RealActivemqQueue;
64
65// NATS: async-nats — 真实实现
66#[cfg(feature = "nats")]
67pub mod real_nats;
68
69#[cfg(feature = "nats")]
70pub use real_nats::RealNatsQueue;
71
72// Pulsar: pulsar crate — 真实实现
73#[cfg(feature = "pulsar")]
74pub mod real_pulsar;
75
76#[cfg(feature = "pulsar")]
77pub use real_pulsar::RealPulsarQueue;
78
79// Kafka: rdkafka — 真实实现
80#[cfg(feature = "kafka")]
81pub mod real_kafka;
82
83#[cfg(feature = "kafka")]
84pub use real_kafka::RealKafkaQueue;
85
86// RocketMQ: 无成熟 Rust 客户端,保持 stub
87// 跟踪项:https://github.com/mxsm/rocketmq-rust (未来可能可用)
88
89#[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        // 验证消息真正进入队列:consume 应能取回,且内容与发布一致
148        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    // ========================================================================
352    // H-3 修复测试:InMemoryQueue 消息大小限制
353    // ========================================================================
354
355    #[tokio::test]
356    async fn test_h3_in_memory_queue_max_messages_limit() {
357        // 设置每 topic 最大 5 条消息
358        let queue = InMemoryQueue::with_max_messages_per_topic(5);
359
360        // 前 5 条成功
361        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        // 第 6 条应失败
370        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        // 消费后可继续发布
377        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        // 默认 100,000 限制,正常使用不应触发
385        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        // 不同 topic 独立计数
398        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        // topic-a 满了
402        assert!(queue.publish("topic-a", b"a3").await.is_err());
403        // topic-b 不受影响
404        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}