Skip to main content

ruststream_amqp/
publisher.rs

1//! [`AmqpPublisher`], its [`AmqpPublish`] policy, and native request/reply.
2
3use std::collections::HashMap;
4use std::sync::Arc;
5use std::time::Duration;
6
7use fe2o3_amqp::Sender;
8use fe2o3_amqp_types::messaging::{Message, Outcome, Properties};
9use ruststream::{OutgoingMessage, PairError, PublishPolicy, Publisher, RequestReply};
10use tokio::sync::Mutex;
11
12use crate::broker::{AmqpCore, ConnectedAmqpBroker, CoreCell};
13use crate::error::{AmqpError, box_err};
14use crate::message::{AmqpMessage, headers_from_amqp, payload_from_body, to_amqp_message};
15
16/// Publishes messages to `AMQP` addresses, one sender link per address, attached lazily on the
17/// shared publisher session.
18///
19/// Buildable before `connect` (it resolves the connection through the broker's shared cell) and
20/// usable until `shutdown`; afterwards every publish reports
21/// [`AmqpError::NotConnected`] instead of silently succeeding against a dead connection.
22#[derive(Clone)]
23pub struct AmqpPublisher {
24    cell: CoreCell,
25    senders: Arc<Mutex<HashMap<String, Arc<Mutex<Sender>>>>>,
26}
27
28impl std::fmt::Debug for AmqpPublisher {
29    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
30        f.debug_struct("AmqpPublisher").finish_non_exhaustive()
31    }
32}
33
34impl AmqpPublisher {
35    pub(crate) fn new(cell: CoreCell) -> Self {
36        Self {
37            cell,
38            senders: Arc::new(Mutex::new(HashMap::new())),
39        }
40    }
41
42    fn core(&self) -> Result<&Arc<AmqpCore>, AmqpError> {
43        let core = self.cell.get().ok_or(AmqpError::NotConnected)?;
44        core.ensure_open()?;
45        Ok(core)
46    }
47
48    /// The sender link for `address`, attached on first use and cached.
49    // The map guard intentionally spans the attach so two callers cannot race a double-attach
50    // for the same address.
51    #[allow(clippy::significant_drop_tightening)]
52    async fn sender_for(
53        &self,
54        core: &AmqpCore,
55        address: &str,
56    ) -> Result<Arc<Mutex<Sender>>, AmqpError> {
57        let mut senders = self.senders.lock().await;
58        if let Some(sender) = senders.get(address) {
59            return Ok(Arc::clone(sender));
60        }
61        let sender = ConnectedAmqpBroker::attach_sender(core, address).await?;
62        let sender = Arc::new(Mutex::new(sender));
63        senders.insert(address.to_owned(), Arc::clone(&sender));
64        Ok(sender)
65    }
66}
67
68/// Sends one built message over a cached sender link and maps a non-accepted outcome to an
69/// error, so a broker-side reject can never pass silently.
70pub(crate) async fn send_message(
71    sender: &Mutex<Sender>,
72    address: &str,
73    message: Message<fe2o3_amqp_types::messaging::Data>,
74) -> Result<(), AmqpError> {
75    let outcome = {
76        let mut sender = sender.lock().await;
77        sender.send(message).await.map_err(|e| AmqpError::Publish {
78            address: address.to_owned(),
79            source: box_err(e),
80        })?
81    };
82    accepted(outcome, address)
83}
84
85pub(crate) fn accepted(outcome: Outcome, address: &str) -> Result<(), AmqpError> {
86    outcome
87        .accepted_or_else(|outcome| AmqpError::PublishNotAccepted {
88            address: address.to_owned(),
89            outcome: format!("{outcome:?}"),
90        })
91        .map(|_| ())
92}
93
94impl Publisher for AmqpPublisher {
95    type Error = AmqpError;
96
97    async fn publish(&self, msg: OutgoingMessage<'_>) -> Result<(), Self::Error> {
98        let core = self.core()?;
99        let sender = self.sender_for(core, msg.name()).await?;
100        send_message(&sender, msg.name(), to_amqp_message(&msg)).await
101    }
102}
103
104impl RequestReply for AmqpPublisher {
105    type Reply = AmqpMessage;
106
107    async fn request(
108        &self,
109        msg: OutgoingMessage<'_>,
110        timeout: Duration,
111    ) -> Result<Self::Reply, Self::Error> {
112        let core = self.core()?;
113
114        // A dynamic receiver per request: the broker names a private reply address that lives
115        // as long as the link. Simple and correct; a shared reply link is a later optimisation.
116        let mut receiver = ConnectedAmqpBroker::attach_dynamic_receiver(core).await?;
117        let reply_to = receiver
118            .source()
119            .as_ref()
120            .and_then(|source| source.address.clone())
121            .ok_or_else(|| AmqpError::Attach {
122                address: "(dynamic)".to_owned(),
123                source: Box::from("the broker did not assign a dynamic reply address"),
124            })?;
125        let correlation_id = core.correlation_id();
126
127        let exchange = async {
128            let mut message = to_amqp_message(&msg);
129            let properties = message.properties.get_or_insert_with(Properties::default);
130            properties.reply_to = Some(reply_to);
131            properties.correlation_id = Some(correlation_id.clone().into());
132
133            let sender = self.sender_for(core, msg.name()).await?;
134            send_message(&sender, msg.name(), message).await?;
135
136            loop {
137                let delivery = receiver
138                    .recv::<fe2o3_amqp_types::messaging::Body<fe2o3_amqp_types::primitives::Value>>(
139                    )
140                    .await
141                    .map_err(|e| AmqpError::Receive {
142                        address: "(dynamic)".to_owned(),
143                        source: box_err(e),
144                    })?;
145                let message = delivery.into_message();
146                let headers = headers_from_amqp(&message);
147                // The private reply address makes foreign traffic unlikely, but correlate
148                // anyway: a late reply to an earlier request must not resolve this one.
149                if headers.correlation_id() == Some(correlation_id.as_str()) {
150                    let payload = payload_from_body(message.body, "(dynamic)")?;
151                    return Ok(AmqpMessage::settled(payload, headers));
152                }
153            }
154        };
155
156        let result = tokio::time::timeout(timeout, exchange)
157            .await
158            .unwrap_or(Err(AmqpError::RequestTimeout));
159        if let Err((_, err)) = receiver.detach().await {
160            tracing::debug!(error = %err, "amqp reply receiver detach failed");
161        }
162        result
163    }
164}
165
166/// The publish policy for [`AmqpPublisher`]: pure declaration, constructible anywhere, paired
167/// with the connected broker by the runtime after `connect`.
168///
169/// # Examples
170///
171/// ```
172/// use ruststream_amqp::AmqpPublish;
173///
174/// let policy = AmqpPublish::default();
175/// # let _ = policy;
176/// ```
177#[derive(Debug, Clone, Copy, Default)]
178#[must_use]
179pub struct AmqpPublish;
180
181impl PublishPolicy<ConnectedAmqpBroker> for AmqpPublish {
182    type Live = AmqpPublisher;
183
184    async fn pair(self, connected: &ConnectedAmqpBroker) -> Result<Self::Live, PairError> {
185        Ok(connected.publisher())
186    }
187}