Skip to main content

rabbit_auto/
auto_ack.rs

1//! Helper to auto ack deliveries
2
3use anyhow::Result;
4use futures::{Stream, StreamExt};
5use lapin::acker::Acker;
6use lapin::message::Delivery;
7use lapin::options::{BasicAckOptions, BasicNackOptions, BasicRejectOptions};
8
9enum Action {
10    Ack(Option<BasicAckOptions>),
11    Nack(Option<BasicNackOptions>),
12    Reject(Option<BasicRejectOptions>),
13}
14/// Automatically ack delivery when this struct drops.
15/// There is an option to automatically nack the delivery.
16pub struct AutoAck {
17    acker: Option<Acker>,
18    action: Action,
19}
20
21impl AutoAck {
22    /// Creates a new auto ack
23    /// # Arguments:
24    /// * channel - the channel where this delivery was received on
25    /// * delivery - delivery which was received. This will take its id of the delivery.
26    pub fn new(acker: Acker) -> Self {
27        Self {
28            acker: Some(acker),
29            action: Action::Ack(None),
30        }
31    }
32
33    /// Creates a new auto ack with BasicAckOptions
34    pub fn new_ack(acker: Acker, options: BasicAckOptions) -> Self {
35        Self {
36            acker: Some(acker),
37            action: Action::Ack(Some(options)),
38        }
39    }
40
41    pub fn new_nack(acker: Acker, options: BasicNackOptions) -> Self {
42        Self {
43            acker: Some(acker),
44            action: Action::Nack(Some(options)),
45        }
46    }
47
48    pub fn new_reject(acker: Acker, options: BasicRejectOptions) -> Self {
49        Self {
50            acker: Some(acker),
51            action: Action::Reject(Some(options)),
52        }
53    }
54
55    /// Change the auto ack in to auto nack
56    pub fn change_to_nack(&mut self, options: Option<BasicNackOptions>) {
57        self.action = Action::Nack(options);
58    }
59
60    pub fn change_to_reject(&mut self, options: Option<BasicRejectOptions>) {
61        self.action = Action::Reject(options);
62    }
63
64    /// Release the channel and the tag from this
65    pub fn release(&mut self) -> Option<Acker> {
66        self.acker.take()
67    }
68
69    /// Perform the ack or nack on the channel, if this has not been already done or released.
70    pub async fn execute(&mut self) -> Result<()> {
71        if let Some(acker) = self.acker.take() {
72            match self.action {
73                Action::Ack(ref mut options) => Self::do_ack(acker, options.take()).await,
74                Action::Nack(ref mut options) => Self::do_nack(acker, options.take()).await,
75                Action::Reject(ref mut option) => Self::do_reject(acker, option.take()).await,
76            }
77        } else {
78            Ok(())
79        }
80    }
81
82    /// Ack the delivery on the channel
83    async fn do_ack(acker: Acker, options: Option<BasicAckOptions>) -> Result<()> {
84        acker
85            .ack(options.unwrap_or_else(|| BasicAckOptions::default()))
86            .await?;
87        Ok(())
88    }
89    /// Nack the delivery on the channel
90    async fn do_nack(acker: Acker, options: Option<BasicNackOptions>) -> Result<()> {
91        acker
92            .nack(options.unwrap_or_else(|| BasicNackOptions::default()))
93            .await?;
94        Ok(())
95    }
96
97    async fn do_reject(acker: Acker, options: Option<BasicRejectOptions>) -> Result<()> {
98        acker
99            .reject(options.unwrap_or_else(|| BasicRejectOptions::default()))
100            .await?;
101        Ok(())
102    }
103}
104
105impl Drop for AutoAck {
106    fn drop(&mut self) {
107        if let Some(acker) = self.acker.take() {
108            match self.action {
109                Action::Ack(ref mut options) => {
110                    tokio::spawn(Self::do_ack(acker, options.take()));
111                }
112                Action::Nack(ref mut options) => {
113                    tokio::spawn(Self::do_nack(acker, options.take()));
114                }
115                Action::Reject(ref mut options) => {
116                    tokio::spawn(Self::do_reject(acker, options.take()));
117                }
118            }
119        }
120    }
121}
122
123/// Map consumer stream into automatically ack stream. The AutoAck object can still be used to release the ack,
124/// and manually ack it or not ack it.
125pub fn auto_ack<S: StreamExt + Stream<Item = Delivery>>(
126    stream: S,
127) -> impl Stream<Item = (AutoAck, Vec<u8>)> {
128    stream.map(|delivery| (AutoAck::new(delivery.acker), delivery.data))
129}