sz-orm-queue 1.2.2

SZ-ORM Message Queue Extension - 6 MQ Providers (RabbitMQ/NATS/Pulsar/Kafka/ActiveMQ real, RocketMQ stub)
Documentation
//! Redis 队列的 InMemory 桩实现(与其它 provider 桩一致,供 `QueueWrapper` 默认使用)。
//!
//! 真实实现见 [`crate::redis_provider::RedisQueueProvider`](启用 `redis` feature)。
//!
//! 说明:本桩模块命名为 `redis_inmemory` 而非 `redis`,是为了避免与外部 `redis` crate
//! 同名冲突(Rust 2018 中本地模块会遮蔽同名外部 crate)。

use crate::error::MqError;
use crate::queue::{InMemoryQueue, Message, MessageQueue};
use async_trait::async_trait;

/// Redis 队列的内存桩实现
///
/// 行为与 [`InMemoryQueue`] 一致,仅作为 `MqProvider::Redis` 在 [`crate::QueueWrapper`]
/// 中的默认承载。真实 Redis 实现请使用 [`crate::RedisQueueProvider`](`redis` feature)。
pub struct InMemoryRedisQueue {
    inner: InMemoryQueue,
}

impl InMemoryRedisQueue {
    /// 创建内存桩实例
    pub fn new() -> Self {
        Self {
            inner: InMemoryQueue::new(),
        }
    }

    /// 当前 topic 的就绪消息数
    pub async fn message_count(&self, topic: &str) -> usize {
        self.inner.message_count(topic).await
    }

    /// 当前 topic 的订阅者数
    pub async fn subscriber_count(&self, topic: &str) -> usize {
        self.inner.subscriber_count(topic).await
    }

    /// 当前 in-flight(已消费未确认)消息数
    pub async fn in_flight_count(&self) -> usize {
        self.inner.in_flight_count().await
    }
}

impl Default for InMemoryRedisQueue {
    fn default() -> Self {
        Self::new()
    }
}

#[async_trait]
impl MessageQueue for InMemoryRedisQueue {
    async fn publish(&self, topic: &str, message: &[u8]) -> Result<(), MqError> {
        self.inner.publish(topic, message).await
    }

    async fn consume(&self, topic: &str) -> Result<Option<Message>, MqError> {
        self.inner.consume(topic).await
    }

    async fn ack(&self, message_id: &str) -> Result<(), MqError> {
        self.inner.ack(message_id).await
    }

    async fn subscribe(&self, topic: &str) -> Result<(), MqError> {
        self.inner.subscribe(topic).await
    }
}

#[cfg(test)]
mod tests {
    use super::*;

    #[tokio::test]
    async fn test_redis_inmemory_publish_and_consume() {
        let queue = InMemoryRedisQueue::new();
        queue.publish("redis-topic", b"hello redis").await.unwrap();
        let msg = queue
            .consume("redis-topic")
            .await
            .unwrap()
            .expect("message should exist");
        assert_eq!(msg.payload, b"hello redis");
        queue.ack(&msg.id).await.unwrap();
    }

    #[tokio::test]
    async fn test_redis_inmemory_consume_empty() {
        let queue = InMemoryRedisQueue::new();
        assert!(queue.consume("empty").await.unwrap().is_none());
    }

    #[tokio::test]
    async fn test_redis_inmemory_subscribe() {
        let queue = InMemoryRedisQueue::new();
        queue.subscribe("redis-topic").await.unwrap();
        assert_eq!(queue.subscriber_count("redis-topic").await, 1);
    }
}