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