1use 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}
14pub struct AutoAck {
17 acker: Option<Acker>,
18 action: Action,
19}
20
21impl AutoAck {
22 pub fn new(acker: Acker) -> Self {
27 Self {
28 acker: Some(acker),
29 action: Action::Ack(None),
30 }
31 }
32
33 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 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 pub fn release(&mut self) -> Option<Acker> {
66 self.acker.take()
67 }
68
69 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 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 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
123pub 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}