sz-orm-queue 1.2.2

SZ-ORM Message Queue Extension - 6 MQ Providers (RabbitMQ/NATS/Pulsar/Kafka/ActiveMQ real, RocketMQ stub)
Documentation
//! RocketMQ provider 集成测试
//!
//! - 非 ignore 测试:trait 实现编译验证、便捷构造(无需真实 RocketMQ)
//! - `#[ignore]` 测试:真实 RocketMQ 5.x Proxy HTTP 端点的发布/消费/确认
//!   运行方式:`cargo test -p sz-orm-queue --features rocketmq -- --ignored`
#![cfg(feature = "rocketmq")]

use sz_orm_queue::queue::MessageQueue;
use sz_orm_queue::RocketMqProvider;

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

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

/// 便捷构造 new:使用默认 topic/group,且 base_url 去尾斜杠
#[tokio::test]
async fn test_rocketmq_provider_new_defaults() {
    let q = RocketMqProvider::new("http://127.0.0.1:8081/")
        .await
        .expect("构造应成功");
    assert_eq!(q.base_url(), "http://127.0.0.1:8081");
    assert_eq!(q.topic(), "default-topic");
    assert_eq!(q.consumer_group(), "default-group");
    // subscribe 为 no-op,始终成功
    q.subscribe("any-topic").await.expect("subscribe 应成功");
}

/// connect:自定义 topic/group,base_url 去尾斜杠
#[tokio::test]
async fn test_rocketmq_provider_connect_custom() {
    let q = RocketMqProvider::connect("http://127.0.0.1:8081/", "my-topic", "my-group")
        .await
        .expect("构造应成功");
    assert_eq!(q.base_url(), "http://127.0.0.1:8081");
    assert_eq!(q.topic(), "my-topic");
    assert_eq!(q.consumer_group(), "my-group");
}

/// 未连接(无服务端)时 publish 应返回错误(而非 panic)
#[tokio::test]
async fn test_rocketmq_provider_publish_no_server_returns_err() {
    let q = RocketMqProvider::new("http://127.0.0.1:1") // 1 号端口几乎不会有服务
        .await
        .expect("构造应成功(不实际连接)");
    let result = q.publish("topic", b"msg").await;
    assert!(result.is_err(), "无服务端时 publish 应失败");
}

// ----------------------------------------------------------------------
// 真实 RocketMQ 集成测试(需 RocketMQ 5.x Proxy HTTP 端点,默认 #[ignore])
// ----------------------------------------------------------------------

/// 完整发布/消费/确认往返
#[tokio::test]
#[ignore = "需真实 RocketMQ 5.x Proxy HTTP 端点"]
async fn integration_rocketmq_round_trip() {
    let queue = RocketMqProvider::connect("http://127.0.0.1:8081", "test-topic", "test-group")
        .await
        .expect("构造应成功");
    queue
        .publish("test-topic", b"hello-rocket")
        .await
        .expect("publish 应成功");
    let msg = queue
        .consume("test-topic")
        .await
        .expect("consume 不应报错")
        .expect("应有消息");
    assert_eq!(msg.payload, b"hello-rocket");
    assert_eq!(msg.topic, "test-topic");
    queue.ack(&msg.id).await.expect("ack 应成功");
}