Skip to main content

ruststream_lapin/
message.rs

1//! The delivery type yielded by [`LapinSubscriber`](crate::LapinSubscriber).
2
3use bytes::Bytes;
4use lapin::Acker;
5use lapin::message::Delivery;
6use lapin::options::{BasicAckOptions, BasicNackOptions, BasicRejectOptions};
7use ruststream::{AckError, Headers, IncomingMessage};
8
9use crate::convert;
10
11/// One AMQP delivery, settled with the protocol's native acknowledgement frames.
12///
13/// Settlement mapping:
14///
15/// - [`ack`](IncomingMessage::ack) sends `basic.ack`.
16/// - [`nack(true)`](IncomingMessage::nack) sends `basic.nack` with `requeue = true`; the broker
17///   redelivers the message (typically to the same queue, `redelivered` set).
18/// - [`nack(false)`](IncomingMessage::nack) sends `basic.reject` with `requeue = false`; the
19///   broker drops the message, or dead-letters it when the queue has a dead-letter exchange.
20///
21/// Replies received through [`LapinRequester`](crate::LapinRequester) arrive on a no-ack
22/// consumer; settling them is a no-op that always succeeds.
23#[derive(Debug)]
24pub struct LapinMessage {
25    payload: Bytes,
26    headers: Headers,
27    exchange: String,
28    routing_key: String,
29    redelivered: bool,
30    delivery_tag: u64,
31    acker: Option<Acker>,
32}
33
34impl LapinMessage {
35    pub(crate) fn from_delivery(delivery: Delivery) -> Self {
36        let headers = convert::headers_from_properties(&delivery.properties);
37        Self {
38            payload: Bytes::from(delivery.data),
39            headers,
40            exchange: delivery.exchange.to_string(),
41            routing_key: delivery.routing_key.to_string(),
42            redelivered: delivery.redelivered,
43            delivery_tag: delivery.delivery_tag,
44            acker: Some(delivery.acker),
45        }
46    }
47
48    /// Builds a settled-by-construction message for no-ack deliveries (request replies).
49    pub(crate) fn from_delivery_no_ack(delivery: Delivery) -> Self {
50        let mut msg = Self::from_delivery(delivery);
51        msg.acker = None;
52        msg
53    }
54
55    /// The exchange this message was published to (empty for the default exchange).
56    #[must_use]
57    pub fn exchange(&self) -> &str {
58        &self.exchange
59    }
60
61    /// The routing key the message was published with.
62    #[must_use]
63    pub fn routing_key(&self) -> &str {
64        &self.routing_key
65    }
66
67    /// Whether the broker marked this delivery as redelivered.
68    #[must_use]
69    pub fn redelivered(&self) -> bool {
70        self.redelivered
71    }
72
73    /// The channel-local delivery tag of this delivery.
74    #[must_use]
75    pub fn delivery_tag(&self) -> u64 {
76        self.delivery_tag
77    }
78
79    async fn settle<F, Fut>(mut self, op: F, what: &'static str) -> Result<(), AckError>
80    where
81        F: FnOnce(Acker) -> Fut,
82        Fut: Future<Output = lapin::Result<bool>>,
83    {
84        // No acker means a no-ack consumer delivered this message; nothing to settle.
85        let Some(acker) = self.acker.take() else {
86            return Ok(());
87        };
88        match op(acker).await {
89            Ok(true) => Ok(()),
90            // lapin reports `false` when the settle frame could not be sent because the channel
91            // already closed or errored; surface that instead of pretending the broker saw it.
92            Ok(false) => Err(AckError::Broker(
93                format!("{what} was not sent: the delivery channel is closed or errored").into(),
94            )),
95            Err(err) => Err(AckError::Broker(Box::new(err))),
96        }
97    }
98}
99
100impl IncomingMessage for LapinMessage {
101    fn payload(&self) -> &[u8] {
102        &self.payload
103    }
104
105    fn headers(&self) -> &Headers {
106        &self.headers
107    }
108
109    /// Acknowledges the delivery with `basic.ack`.
110    ///
111    /// # Errors
112    ///
113    /// Returns [`AckError::Broker`] when the frame cannot be sent, for example because the
114    /// channel closed after the delivery arrived.
115    ///
116    /// # Cancel safety
117    ///
118    /// Not cancel safe: dropping the future after the frame was queued may still acknowledge the
119    /// message on the broker.
120    async fn ack(self) -> Result<(), AckError> {
121        self.settle(
122            |acker| async move { acker.ack(BasicAckOptions::default()).await },
123            "basic.ack",
124        )
125        .await
126    }
127
128    /// Settles negatively: `basic.nack(requeue = true)` or `basic.reject(requeue = false)`.
129    ///
130    /// # Errors
131    ///
132    /// Returns [`AckError::Broker`] when the frame cannot be sent, for example because the
133    /// channel closed after the delivery arrived.
134    ///
135    /// # Cancel safety
136    ///
137    /// Not cancel safe: dropping the future after the frame was queued may still settle the
138    /// message on the broker.
139    async fn nack(self, requeue: bool) -> Result<(), AckError> {
140        if requeue {
141            self.settle(
142                |acker| async move {
143                    acker
144                        .nack(BasicNackOptions {
145                            multiple: false,
146                            requeue: true,
147                        })
148                        .await
149                },
150                "basic.nack",
151            )
152            .await
153        } else {
154            self.settle(
155                |acker| async move { acker.reject(BasicRejectOptions { requeue: false }).await },
156                "basic.reject",
157            )
158            .await
159        }
160    }
161}