sz-orm-queue 1.2.2

SZ-ORM Message Queue Extension - 6 MQ Providers (RabbitMQ/NATS/Pulsar/Kafka/ActiveMQ real, RocketMQ stub)
Documentation
//! Redis provider 集成测试
//!
//! - 非 ignore 测试:配置构造、trait 实现编译验证(无需真实 Redis)
//! - `#[ignore]` 测试:真实 Redis 连接的发布/消费/确认
//!   启动方式:`docker run -d -p 6379:6379 redis:7`
//!   运行方式:`cargo test -p sz-orm-queue --features redis -- --ignored`
#![cfg(feature = "redis")]

use sz_orm_queue::queue::{MessageQueue, RedisConfig, RedisMode};
use sz_orm_queue::RedisQueueProvider;

// ----------------------------------------------------------------------
// 单元测试(无需真实服务)
// ----------------------------------------------------------------------

/// 默认配置:List 模式、无 URL、无消费者组
#[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);
}

/// Stream 模式配置可正确构造
#[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"));
}

/// 三种 RedisMode 互不相等
#[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);
}

/// 编译验证:RedisQueueProvider 实现了 MessageQueue trait
#[test]
fn test_redis_provider_implements_message_queue() {
    fn assert_message_queue<T: MessageQueue>() {}
    assert_message_queue::<RedisQueueProvider>();
}

// ----------------------------------------------------------------------
// 真实 Redis 集成测试(需启动 Redis,默认 #[ignore])
// ----------------------------------------------------------------------

/// List 模式:RPUSH + BLPOP,FIFO,消费即出队,ack 为 no-op
#[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);
    // List 模式 ack 为 no-op,始终成功
    queue.ack(&msg.id).await.expect("ack 应成功");
}

/// Stream 模式:XADD + XREADGROUP + XACK,消费者组,可靠投递
#[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");
    // Stream 模式 key 应为 entry id(非空)
    assert!(msg.key.is_some(), "Stream 消息应带 entry id");
    // message_id 形如 "{topic}::{entry_id}"
    assert!(msg.id.contains("::"), "message_id 应含 :: 分隔符");
    // Stream 模式 ack 真正调用 XACK
    queue.ack(&msg.id).await.expect("ack 应成功");
}

/// PubSub 模式:PUBLISH + SUBSCRIBE,需先订阅再发布
#[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");
}