ruststream_amqp/
publisher.rs1use 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#[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 #[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
68pub(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 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 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#[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}