#![cfg(feature = "redis")]
use sz_orm_queue::queue::{MessageQueue, RedisConfig, RedisMode};
use sz_orm_queue::RedisQueueProvider;
#[test]
fn test_redis_provider_config_default() {
let cfg = RedisConfig::default();
assert_eq!(cfg.mode, RedisMode::List);
assert!(cfg.url.is_none());
assert!(cfg.consumer_group.is_none());
assert!(cfg.consumer_name.is_none());
assert_eq!(cfg.pool_size, 8);
}
#[test]
fn test_redis_provider_stream_config_construction() {
let cfg = RedisConfig {
url: Some("redis://127.0.0.1:6379/0".into()),
mode: RedisMode::Stream,
consumer_group: Some("test-group".into()),
consumer_name: Some("test-consumer".into()),
pool_size: 4,
};
assert_eq!(cfg.mode, RedisMode::Stream);
assert_eq!(cfg.consumer_group.as_deref(), Some("test-group"));
assert_eq!(cfg.consumer_name.as_deref(), Some("test-consumer"));
}
#[test]
fn test_redis_mode_distinct_variants() {
assert_ne!(RedisMode::List, RedisMode::PubSub);
assert_ne!(RedisMode::PubSub, RedisMode::Stream);
assert_ne!(RedisMode::List, RedisMode::Stream);
}
#[test]
fn test_redis_provider_implements_message_queue() {
fn assert_message_queue<T: MessageQueue>() {}
assert_message_queue::<RedisQueueProvider>();
}
#[tokio::test]
#[ignore = "需真实 Redis 服务器"]
async fn integration_redis_list_round_trip() {
let queue = RedisQueueProvider::new("redis://127.0.0.1:6379/0")
.await
.expect("Redis 连接应成功");
let topic = "sz-orm-queue-test-list";
queue
.publish(topic, b"hello-list")
.await
.expect("publish 应成功");
let msg = queue
.consume(topic)
.await
.expect("consume 不应报错")
.expect("应有消息");
assert_eq!(msg.payload, b"hello-list");
assert_eq!(msg.topic, topic);
queue.ack(&msg.id).await.expect("ack 应成功");
}
#[tokio::test]
#[ignore = "需真实 Redis 服务器"]
async fn integration_redis_stream_round_trip() {
let queue = RedisQueueProvider::connect(RedisConfig {
url: Some("redis://127.0.0.1:6379/0".into()),
mode: RedisMode::Stream,
..Default::default()
})
.await
.expect("Redis 连接应成功");
let topic = "sz-orm-queue-test-stream";
queue
.publish(topic, b"hello-stream")
.await
.expect("publish 应成功");
let msg = queue
.consume(topic)
.await
.expect("consume 不应报错")
.expect("应有消息");
assert_eq!(msg.payload, b"hello-stream");
assert!(msg.key.is_some(), "Stream 消息应带 entry id");
assert!(msg.id.contains("::"), "message_id 应含 :: 分隔符");
queue.ack(&msg.id).await.expect("ack 应成功");
}
#[tokio::test]
#[ignore = "需真实 Redis 服务器"]
async fn integration_redis_pubsub_round_trip() {
let queue = RedisQueueProvider::connect(RedisConfig {
url: Some("redis://127.0.0.1:6379/0".into()),
mode: RedisMode::PubSub,
..Default::default()
})
.await
.expect("Redis 连接应成功");
let topic = "sz-orm-queue-test-pubsub";
queue.subscribe(topic).await.expect("subscribe 应成功");
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
queue
.publish(topic, b"hello-pubsub")
.await
.expect("publish 应成功");
let msg = queue
.consume(topic)
.await
.expect("consume 不应报错")
.expect("应有消息");
assert_eq!(msg.payload, b"hello-pubsub");
}