use fe2o3_amqp_types::primitives::OrderedMap;
use serde_amqp::Value;
use crate::amqp::amqp_receiver::AmqpReceiver;
use crate::primitives::service_bus_peeked_message::ServiceBusPeekedMessage;
use crate::util::IntoAzureCoreError;
use crate::{
core::TransportReceiver, primitives::service_bus_received_message::ServiceBusReceivedMessage,
};
use crate::{primitives::sub_queue::SubQueue, ServiceBusReceiveMode};
use super::DeadLetterOptions;
#[cfg(docsrs)]
use crate::{ServiceBusClient, ServiceBusRetryOptions};
#[derive(Debug, Default, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub struct ServiceBusReceiverOptions {
pub prefetch_count: u32,
pub receive_mode: ServiceBusReceiveMode,
pub identifier: Option<String>,
pub sub_queue: SubQueue,
}
#[derive(Debug)]
pub struct ServiceBusReceiver {
pub(crate) inner: AmqpReceiver,
}
impl ServiceBusReceiver {
pub fn entity_path(&self) -> &str {
self.inner.entity_path()
}
pub fn identifier(&self) -> &str {
self.inner.identifier()
}
pub fn prefetch_count(&self) -> u32 {
self.inner.prefetch_count()
}
pub fn receive_mode(&self) -> ServiceBusReceiveMode {
self.inner.receive_mode()
}
pub async fn dispose(self) -> Result<(), azure_core::Error> {
self.inner.close().await.map_err(IntoAzureCoreError::into_azure_core_error)
}
pub async fn receive_message(
&mut self,
) -> Result<ServiceBusReceivedMessage, azure_core::Error> {
self.receive_messages(1).await.map(|mut v| {
v.drain(..)
.next()
.expect("At least one message should be received.")
})
}
pub async fn receive_messages(
&mut self,
max_messages: u32,
) -> Result<Vec<ServiceBusReceivedMessage>, azure_core::Error> {
self.inner.receive_messages(max_messages).await.map_err(Into::into)
}
pub async fn receive_message_with_max_wait_time(
&mut self,
max_wait_time: impl Into<Option<std::time::Duration>>,
) -> Result<Option<ServiceBusReceivedMessage>, azure_core::Error> {
self.receive_messages_with_max_wait_time(1, max_wait_time)
.await
.map(|mut v| v.drain(..).next())
}
pub async fn receive_messages_with_max_wait_time(
&mut self,
max_messages: u32,
max_wait_time: impl Into<Option<std::time::Duration>>,
) -> Result<Vec<ServiceBusReceivedMessage>, azure_core::Error> {
self.inner
.receive_messages_with_max_wait_time(max_messages, max_wait_time.into())
.await
.map_err(Into::into)
}
pub async fn complete_message(
&mut self,
message: impl AsRef<ServiceBusReceivedMessage>,
) -> Result<(), azure_core::Error> {
self.inner.complete(message.as_ref(), None).await.map_err(Into::into)
}
pub async fn abandon_message(
&mut self,
message: impl AsRef<ServiceBusReceivedMessage>,
properties_to_modify: Option<OrderedMap<String, Value>>,
) -> Result<(), azure_core::Error> {
self.inner
.abandon(message.as_ref(), properties_to_modify, None)
.await
.map_err(Into::into)
}
pub async fn dead_letter_message(
&mut self,
message: impl AsRef<ServiceBusReceivedMessage>,
options: DeadLetterOptions,
) -> Result<(), azure_core::Error> {
self.inner
.dead_letter(
message.as_ref(),
options.dead_letter_reason,
options.dead_letter_error_description,
options.properties_to_modify,
None,
)
.await
.map_err(Into::into)
}
pub async fn defer_message(
&mut self,
message: impl AsRef<ServiceBusReceivedMessage>,
properties_to_modify: Option<OrderedMap<String, Value>>,
) -> Result<(), azure_core::Error> {
self.inner
.defer(message.as_ref(), properties_to_modify, None)
.await
.map_err(Into::into)
}
pub async fn peek_message(
&mut self,
from_sequence_number: Option<i64>,
) -> Result<Option<ServiceBusPeekedMessage>, azure_core::Error> {
self.peek_messages(1, from_sequence_number)
.await
.map(|mut v| v.drain(..).next())
}
pub async fn peek_messages(
&mut self,
max_messages: u32, from_sequence_number: Option<i64>,
) -> Result<Vec<ServiceBusPeekedMessage>, azure_core::Error> {
self.inner
.peek_messages(from_sequence_number, max_messages as i32)
.await
.map_err(Into::into)
}
pub async fn receive_deferred_message(
&mut self,
sequence_number: i64,
) -> Result<Option<ServiceBusReceivedMessage>, azure_core::Error> {
self.receive_deferred_messages(std::iter::once(sequence_number))
.await
.map(|mut v| v.drain(..).next())
}
pub async fn receive_deferred_messages<Seq>(
&mut self,
sequence_numbers: Seq,
) -> Result<Vec<ServiceBusReceivedMessage>, azure_core::Error>
where
Seq: IntoIterator<Item = i64> + Send,
Seq::IntoIter: Send,
{
self.inner
.receive_deferred_messages(sequence_numbers.into_iter(), None)
.await
.map_err(Into::into)
}
pub async fn renew_message_lock(
&mut self,
message: &mut ServiceBusReceivedMessage,
) -> Result<(), azure_core::Error> {
let lock_tokens = vec![message.lock_token().clone()];
let mut expirations = self.inner.renew_message_lock(lock_tokens).await?;
if let Some(expiration) = expirations.drain(..).next() {
message.set_locked_until(expiration);
}
Ok(())
}
}