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