#![cfg(feature = "rocketmq")]
use sz_orm_queue::queue::MessageQueue;
use sz_orm_queue::RocketMqProvider;
#[test]
fn test_rocketmq_provider_implements_message_queue() {
fn assert_message_queue<T: MessageQueue>() {}
assert_message_queue::<RocketMqProvider>();
}
#[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");
q.subscribe("any-topic").await.expect("subscribe 应成功");
}
#[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");
}
#[tokio::test]
async fn test_rocketmq_provider_publish_no_server_returns_err() {
let q = RocketMqProvider::new("http://127.0.0.1:1") .await
.expect("构造应成功(不实际连接)");
let result = q.publish("topic", b"msg").await;
assert!(result.is_err(), "无服务端时 publish 应失败");
}
#[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 应成功");
}