pulses 0.2.0

A robust, high-performance background job processing library for Rust.
Documentation
//! Periodically reclaims pending messages abandoned by crashed/stalled
//! consumers and re-routes them to handler pools.

use std::time::Duration;

use tokio::time::MissedTickBehavior;
use tokio_util::sync::CancellationToken;

use crate::core::Broker;
use crate::runtime::router::Router;

pub(crate) struct Reclaimer<B: Broker> {
    broker: B,
    subscription: B::Subscription,
    router: Router<B::AckToken>,
    interval: Duration,
    cancellation: CancellationToken,
}

impl<B: Broker> Reclaimer<B> {
    pub(crate) fn new(
        broker: B, subscription: B::Subscription, router: Router<B::AckToken>, interval: Duration,
        cancellation: CancellationToken,
    ) -> Self {
        Self { broker, subscription, router, interval, cancellation }
    }

    pub(crate) async fn run(self) {
        let mut ticker = tokio::time::interval(self.interval.max(Duration::from_millis(1)));
        ticker.set_missed_tick_behavior(MissedTickBehavior::Delay);

        loop {
            tokio::select! {
                _ = self.cancellation.cancelled() => return,
                _ = ticker.tick() => {}
            }

            let reclaim_result = tokio::select! {
                _ = self.cancellation.cancelled() => return,
                result = self.broker.reclaim(&self.subscription) => result,
            };

            match reclaim_result {
                Ok(batch) => {
                    for (envelope, token) in batch {
                        if let Err(error) = self.router.route(envelope, token, &self.cancellation).await {
                            tracing::error!(%error, "failed to route reclaimed message");
                        }
                        if self.cancellation.is_cancelled() {
                            return;
                        }
                    }
                }
                Err(error) => {
                    tracing::warn!(%error, "broker reclaim failed");
                }
            }
        }
    }
}

#[cfg(test)]
mod tests {
    use std::sync::Arc;
    use std::time::Duration;

    use tokio::sync::mpsc;
    use tokio::time::timeout;
    use tokio_util::sync::CancellationToken;

    use super::Reclaimer;
    use crate::handler_set::RoutingTable;
    use crate::runtime::router::Router;
    use crate::test_support::MockBroker;
    use crate::test_support::MockSubscription;
    use crate::test_support::make_envelope;
    use crate::test_support::token;

    #[tokio::test]
    async fn reclaims_and_routes() {
        let broker = MockBroker::default();
        broker.enqueue_reclaim(vec![(make_envelope("orders", "9-9"), token("9-9"))]);

        let (tx, mut rx) = mpsc::channel(8);
        let mut table = RoutingTable::new();
        table.insert(Arc::from("orders"), 1 << 0);
        let router = Router::new(table, vec![tx]);
        let cancellation = CancellationToken::new();

        let reclaimer =
            Reclaimer::new(broker, MockSubscription::default(), router, Duration::from_millis(5), cancellation.clone());
        let handle = tokio::spawn(reclaimer.run());

        let task = timeout(Duration::from_secs(1), rx.recv()).await.unwrap().unwrap();
        assert_eq!(task.envelope.id.as_ref(), "9-9");

        cancellation.cancel();
        timeout(Duration::from_secs(1), handle).await.unwrap().unwrap();
    }
}