use std::{borrow::Cow, marker::PhantomData, sync::Arc};
use azure_core::{credentials::TokenCredential, http::Url};
use crate::{
amqp::{
amqp_client::AmqpClient,
error::OpenReceiverError,
},
authorization::{
service_bus_token_credential::ServiceBusTokenCredential,
shared_access_credential::SharedAccessCredential, AzureNamedKeyCredential,
AzureSasCredential,
},
core::{BasicRetryPolicy, TransportSessionReceiver},
diagnostics,
entity_name_formatter::{self, format_entity_path},
primitives::{
service_bus_connection::{build_connection_resource, ServiceBusConnection},
service_bus_retry_options::ServiceBusRetryOptions,
service_bus_retry_policy::ServiceBusRetryPolicyExt,
service_bus_transport_type::ServiceBusTransportType,
},
receiver::service_bus_session_receiver::{
ServiceBusSessionReceiver, ServiceBusSessionReceiverOptions,
},
ServiceBusReceiver, ServiceBusReceiverOptions, ServiceBusRuleManager, ServiceBusSender,
ServiceBusSenderOptions,
};
use super::error::AcceptNextSessionError;
#[derive(Debug, Clone, Default)]
pub struct ServiceBusClientOptions {
pub transport_type: ServiceBusTransportType,
pub identifier: Option<String>,
pub custom_endpoint_address: Option<Url>,
pub retry_options: ServiceBusRetryOptions,
pub enable_cross_entity_transactions: bool,
}
#[derive(Debug)]
pub struct WithCustomRetryPolicy<RP> {
retry_policy: PhantomData<RP>,
}
impl<RP> WithCustomRetryPolicy<RP>
where
RP: ServiceBusRetryPolicyExt + Send + Sync + 'static,
{
pub async fn new_from_connection_string<'a>(
self,
connection_string: impl Into<Cow<'a, str>>,
options: ServiceBusClientOptions,
) -> Result<ServiceBusClient<RP>, azure_core::Error> {
let connection_string = connection_string.into();
let identifier = options.identifier.clone();
let connection = ServiceBusConnection::new(connection_string, options).await?;
let identifier = identifier.unwrap_or_else(|| {
diagnostics::utilities::generate_identifier(connection.fully_qualified_namespace())
});
Ok(ServiceBusClient {
identifier,
connection,
})
}
#[deprecated(
since = "0.14.0",
note = "Please use `new_from_connection_string` instead"
)]
pub async fn create_client<'a>(
self,
connection_string: impl Into<Cow<'a, str>>,
options: ServiceBusClientOptions,
) -> Result<ServiceBusClient<RP>, azure_core::Error> {
self.new_from_connection_string(connection_string, options).await
}
pub async fn new_from_named_key_credential(
self,
fully_qualified_namespace: impl Into<String>,
credential: AzureNamedKeyCredential,
options: ServiceBusClientOptions,
) -> Result<ServiceBusClient<RP>, azure_core::Error> {
let fully_qualified_namespace = fully_qualified_namespace.into();
let signuture_resource = build_connection_resource(
&options.transport_type,
Some(&fully_qualified_namespace),
None,
)?;
let shared_access_credential =
SharedAccessCredential::try_from_named_key_credential(credential, signuture_resource)?;
self.new_from_credential(fully_qualified_namespace, shared_access_credential, options).await
}
#[deprecated(
since = "0.14.0",
note = "Please use `new_from_named_key_credential` instead"
)]
pub async fn create_client_with_named_key_credential(
self,
fully_qualified_namespace: impl Into<String>,
credential: AzureNamedKeyCredential,
options: ServiceBusClientOptions,
) -> Result<ServiceBusClient<RP>, azure_core::Error> {
self.new_from_named_key_credential(fully_qualified_namespace, credential, options).await
}
pub async fn new_from_sas_credential(
self,
fully_qualified_namespace: impl Into<String>,
credential: AzureSasCredential,
options: ServiceBusClientOptions,
) -> Result<ServiceBusClient<RP>, azure_core::Error> {
let shared_access_credential = SharedAccessCredential::try_from_sas_credential(credential)?;
self.new_from_credential(fully_qualified_namespace, shared_access_credential, options).await
}
#[deprecated(
since = "0.14.0",
note = "Please use `new_from_sas_credential` instead"
)]
pub async fn create_client_with_sas_credential(
self,
fully_qualified_namespace: impl Into<String>,
credential: AzureSasCredential,
options: ServiceBusClientOptions,
) -> Result<ServiceBusClient<RP>, azure_core::Error> {
self.new_from_sas_credential(fully_qualified_namespace, credential, options).await
}
pub async fn new_from_token_credential(
self,
fully_qualified_namespace: impl Into<String>,
credential: Arc<dyn TokenCredential>,
options: ServiceBusClientOptions,
) -> Result<ServiceBusClient<RP>, azure_core::Error> {
let credential = ServiceBusTokenCredential::new(credential);
self.new_from_credential(fully_qualified_namespace, credential, options).await
}
#[deprecated(
since = "0.14.0",
note = "Please use `new_from_token_credential` instead"
)]
pub async fn create_client_with_token_credential(
self,
fully_qualified_namespace: impl Into<String>,
credential: Arc<dyn TokenCredential>,
options: ServiceBusClientOptions,
) -> Result<ServiceBusClient<RP>, azure_core::Error> {
self.new_from_token_credential(fully_qualified_namespace, credential, options).await
}
pub async fn new_from_credential(
self,
fully_qualified_namespace: impl Into<String>,
credential: impl Into<ServiceBusTokenCredential>,
options: ServiceBusClientOptions,
) -> Result<ServiceBusClient<RP>, azure_core::Error> {
let fully_qualified_namespace = fully_qualified_namespace.into();
let identifier = options.identifier.clone().unwrap_or_else(|| {
diagnostics::utilities::generate_identifier(&fully_qualified_namespace)
});
let credential = credential.into();
let connection = ServiceBusConnection::new_from_credential(
fully_qualified_namespace.into(),
credential,
options,
)
.await?;
Ok(ServiceBusClient {
identifier,
connection,
})
}
}
cfg_unsecured! {
impl<RP> WithCustomRetryPolicy<RP> {
pub fn unsecured() -> Unsecured<RP> {
Unsecured {
retry_policy: PhantomData,
}
}
}
#[derive(Debug)]
pub struct Unsecured<RP> {
retry_policy: PhantomData<RP>,
}
impl<RP> Unsecured<RP>
where
RP: ServiceBusRetryPolicyExt + Send + Sync + 'static,
{
pub async fn new_from_connection_string<'a>(
self,
connection_string: impl Into<Cow<'a, str>>,
options: ServiceBusClientOptions,
) -> Result<ServiceBusClient<RP>, azure_core::Error> {
let connection_string = connection_string.into();
let identifier = options.identifier.clone();
let connection = ServiceBusConnection::new_unsecured(connection_string, options).await?;
let identifier = identifier.unwrap_or_else(|| {
diagnostics::utilities::generate_identifier(connection.fully_qualified_namespace())
});
Ok(ServiceBusClient {
identifier,
connection,
})
}
pub async fn new_from_named_key_credential(
self,
fully_qualified_namespace: impl Into<String>,
credential: AzureNamedKeyCredential,
options: ServiceBusClientOptions,
) -> Result<ServiceBusClient<RP>, azure_core::Error> {
let fully_qualified_namespace = fully_qualified_namespace.into();
let signuture_resource = crate::primitives::service_bus_connection::build_unsecured_connection_resource(
&options.transport_type,
Some(&fully_qualified_namespace),
None,
)?;
let shared_access_credential =
SharedAccessCredential::try_from_named_key_credential(credential, signuture_resource)?;
self.new_from_credential(fully_qualified_namespace, shared_access_credential, options).await
}
pub async fn new_from_sas_credential(
self,
fully_qualified_namespace: impl Into<String>,
credential: AzureSasCredential,
options: ServiceBusClientOptions,
) -> Result<ServiceBusClient<RP>, azure_core::Error> {
let shared_access_credential = SharedAccessCredential::try_from_sas_credential(credential)?;
self.new_from_credential(fully_qualified_namespace, shared_access_credential, options).await
}
pub async fn new_from_token_credential(
self,
fully_qualified_namespace: impl Into<String>,
credential: Arc<dyn TokenCredential>,
options: ServiceBusClientOptions,
) -> Result<ServiceBusClient<RP>, azure_core::Error> {
let credential = ServiceBusTokenCredential::new(credential);
self.new_from_credential(fully_qualified_namespace, credential, options)
.await
}
pub async fn new_from_credential(
self,
fully_qualified_namespace: impl Into<String>,
credential: impl Into<ServiceBusTokenCredential>,
options: ServiceBusClientOptions,
) -> Result<ServiceBusClient<RP>, azure_core::Error> {
let fully_qualified_namespace = fully_qualified_namespace.into();
let identifier = options.identifier.clone().unwrap_or_else(|| {
diagnostics::utilities::generate_identifier(&fully_qualified_namespace)
});
let credential = credential.into();
let connection = ServiceBusConnection::new_unsecured_from_credential(
fully_qualified_namespace,
credential,
options,
)
.await?;
Ok(ServiceBusClient {
identifier,
connection,
})
}
}
}
#[derive(Debug)]
pub struct ServiceBusClient<RP> {
identifier: String,
connection: ServiceBusConnection<AmqpClient<RP>>,
}
impl ServiceBusClient<BasicRetryPolicy> {
pub fn with_custom_retry_policy<RP>() -> WithCustomRetryPolicy<RP> {
WithCustomRetryPolicy {
retry_policy: PhantomData,
}
}
cfg_unsecured! {
pub fn unsecured() -> Unsecured<BasicRetryPolicy> {
Unsecured {
retry_policy: PhantomData,
}
}
}
pub async fn new_from_connection_string<'a>(
connection_string: impl Into<Cow<'a, str>>,
options: ServiceBusClientOptions,
) -> Result<Self, azure_core::Error> {
Self::with_custom_retry_policy()
.new_from_connection_string(connection_string, options)
.await
}
#[deprecated(
since = "0.14.0",
note = "Please use `new_from_connection_string` instead"
)]
pub async fn new<'a>(
connection_string: impl Into<Cow<'a, str>>,
options: ServiceBusClientOptions,
) -> Result<Self, azure_core::Error> {
Self::with_custom_retry_policy()
.new_from_connection_string(connection_string, options)
.await
}
pub async fn new_from_named_key_credential(
fully_qualified_namespace: impl Into<String>,
credential: AzureNamedKeyCredential,
options: ServiceBusClientOptions,
) -> Result<Self, azure_core::Error> {
Self::with_custom_retry_policy()
.new_from_named_key_credential(fully_qualified_namespace, credential, options)
.await
}
#[deprecated(
since = "0.14.0",
note = "Please use `new_from_named_key_credential` instead"
)]
pub async fn new_with_named_key_credential(
fully_qualified_namespace: impl Into<String>,
credential: AzureNamedKeyCredential,
options: ServiceBusClientOptions,
) -> Result<Self, azure_core::Error> {
Self::new_from_named_key_credential(fully_qualified_namespace, credential, options).await
}
pub async fn new_from_sas_credential(
fully_qualified_namespace: impl Into<String>,
credential: AzureSasCredential,
options: ServiceBusClientOptions,
) -> Result<Self, azure_core::Error> {
Self::with_custom_retry_policy()
.new_from_sas_credential(fully_qualified_namespace, credential, options)
.await
}
#[deprecated(
since = "0.14.0",
note = "Please use `new_from_sas_credential` instead"
)]
pub async fn new_with_sas_credential(
fully_qualified_namespace: impl Into<String>,
credential: AzureSasCredential,
options: ServiceBusClientOptions,
) -> Result<Self, azure_core::Error> {
Self::new_from_sas_credential(fully_qualified_namespace, credential, options).await
}
pub async fn new_from_token_credential(
fully_qualified_namespace: impl Into<String>,
credential: Arc<dyn TokenCredential>,
options: ServiceBusClientOptions,
) -> Result<Self, azure_core::Error> {
Self::with_custom_retry_policy()
.new_from_token_credential(fully_qualified_namespace, credential, options)
.await
}
#[deprecated(
since = "0.14.0",
note = "Please use `new_from_token_credential` instead"
)]
pub async fn new_with_token_credential(
fully_qualified_namespace: impl Into<String>,
credential: Arc<dyn TokenCredential>,
options: ServiceBusClientOptions,
) -> Result<Self, azure_core::Error> {
Self::new_from_token_credential(fully_qualified_namespace, credential, options).await
}
pub async fn new_from_credential(
fully_qualified_namespace: impl Into<String>,
credential: impl Into<ServiceBusTokenCredential>,
options: ServiceBusClientOptions,
) -> Result<Self, azure_core::Error> {
Self::with_custom_retry_policy()
.new_from_credential(fully_qualified_namespace, credential, options)
.await
}
}
impl<RP> ServiceBusClient<RP>
where
RP: ServiceBusRetryPolicyExt + 'static,
{
pub fn fully_qualified_namespace(&self) -> &str {
self.connection.fully_qualified_namespace()
}
pub fn identifier(&self) -> &str {
&self.identifier
}
pub fn is_closed(&self) -> bool {
self.connection.is_closed()
}
}
impl<RP> ServiceBusClient<RP>
where
RP: ServiceBusRetryPolicyExt + 'static,
{
pub async fn dispose(self) -> Result<(), azure_core::Error> {
self.connection.dispose().await?;
Ok(())
}
}
impl<RP> ServiceBusClient<RP>
where
RP: ServiceBusRetryPolicyExt + 'static,
{
pub async fn create_sender(
&mut self,
queue_or_topic_name: impl Into<String>,
options: ServiceBusSenderOptions,
) -> Result<ServiceBusSender, azure_core::Error> {
let entity_path = queue_or_topic_name.into();
let identifier = options
.identifier
.filter(|id| !id.is_empty())
.unwrap_or_else(|| diagnostics::utilities::generate_identifier(&entity_path));
let retry_options = self.connection.retry_options().clone();
let inner = self
.connection
.create_transport_sender(entity_path.into(), identifier, retry_options)
.await?;
Ok(ServiceBusSender { inner })
}
}
impl<RP> ServiceBusClient<RP>
where
RP: ServiceBusRetryPolicyExt + 'static,
{
pub fn transport_type(&self) -> ServiceBusTransportType {
self.connection.transport_type()
}
pub async fn create_receiver_for_queue(
&mut self,
queue_name: impl Into<String>,
options: ServiceBusReceiverOptions,
) -> Result<ServiceBusReceiver, azure_core::Error> {
let entity_path = queue_name.into();
self.create_receiver(entity_path, options).await
.map_err(Into::into)
}
pub async fn create_receiver_for_subscription(
&mut self,
topic_name: impl AsRef<str>,
subscription_name: impl AsRef<str>,
options: ServiceBusReceiverOptions,
) -> Result<ServiceBusReceiver, azure_core::Error> {
let entity_path = entity_name_formatter::format_subscription_path(
topic_name.as_ref(),
subscription_name.as_ref(),
);
self.create_receiver(entity_path, options).await
.map_err(Into::into)
}
async fn create_receiver(
&mut self,
entity_path: String,
options: ServiceBusReceiverOptions,
) -> Result<ServiceBusReceiver, OpenReceiverError> {
let identifier = options
.identifier
.filter(|id| !id.is_empty())
.unwrap_or_else(|| diagnostics::utilities::generate_identifier(&entity_path));
let retry_options = self.connection.retry_options().clone();
let receive_mode = options.receive_mode;
let prefetch_count = options.prefetch_count;
let entity_path = format_entity_path(entity_path, options.sub_queue);
let inner = self
.connection
.create_transport_receiver(
entity_path,
identifier,
retry_options,
receive_mode,
prefetch_count,
)
.await?;
Ok(ServiceBusReceiver { inner })
}
pub async fn accept_session_for_queue(
&mut self,
queue_name: impl Into<String>,
session_id: impl Into<String>,
options: ServiceBusSessionReceiverOptions,
) -> Result<ServiceBusSessionReceiver, azure_core::Error> {
let entity_path = queue_name.into();
let session_id = session_id.into();
self.accept_session(entity_path, session_id, options).await.map_err(Into::into)
}
pub async fn accept_session_for_subscription(
&mut self,
topic_name: impl AsRef<str>,
subscription_name: impl AsRef<str>,
session_id: impl Into<String>,
options: ServiceBusSessionReceiverOptions,
) -> Result<ServiceBusSessionReceiver, azure_core::Error> {
let entity_path = entity_name_formatter::format_subscription_path(
topic_name.as_ref(),
subscription_name.as_ref(),
);
let session_id = session_id.into();
self.accept_session(entity_path, session_id, options).await.map_err(Into::into)
}
async fn accept_session(
&mut self,
entity_path: String,
session_id: String,
options: ServiceBusSessionReceiverOptions,
) -> Result<ServiceBusSessionReceiver, OpenReceiverError> {
let identifier = options
.identifier
.unwrap_or_else(|| diagnostics::utilities::generate_identifier(&entity_path));
let retry_options = self.connection.retry_options().clone();
let receive_mode = options.receive_mode;
let prefetch_count = options.prefetch_count;
let inner = self
.connection
.create_transport_session_receiver(
entity_path,
identifier,
retry_options,
receive_mode,
prefetch_count,
Some(session_id.clone()),
)
.await?;
Ok(ServiceBusSessionReceiver { inner, session_id })
}
}
impl<RP> ServiceBusClient<RP>
where
RP: ServiceBusRetryPolicyExt + 'static,
{
pub async fn accept_next_session_for_queue(
&mut self,
queue_name: impl Into<String>,
options: ServiceBusSessionReceiverOptions,
) -> Result<ServiceBusSessionReceiver, azure_core::Error> {
let entity_path = queue_name.into();
self.accept_next_session(entity_path, options).await.map_err(Into::into)
}
pub async fn accept_next_session_for_subscription(
&mut self,
topic_name: impl AsRef<str>,
subscription_name: impl AsRef<str>,
options: ServiceBusSessionReceiverOptions,
) -> Result<ServiceBusSessionReceiver, azure_core::Error> {
let entity_path = entity_name_formatter::format_subscription_path(
topic_name.as_ref(),
subscription_name.as_ref(),
);
self.accept_next_session(entity_path, options).await.map_err(Into::into)
}
async fn accept_next_session(
&mut self,
entity_path: String,
options: ServiceBusSessionReceiverOptions,
) -> Result<ServiceBusSessionReceiver, AcceptNextSessionError> {
let identifier = options
.identifier
.unwrap_or_else(|| diagnostics::utilities::generate_identifier(&entity_path));
let retry_options = self.connection.retry_options().clone();
let receive_mode = options.receive_mode;
let prefetch_count = options.prefetch_count;
let inner = self
.connection
.create_transport_session_receiver(
entity_path,
identifier,
retry_options,
receive_mode,
prefetch_count,
None,
)
.await?;
let session_id = inner
.session_id()
.ok_or(AcceptNextSessionError::SessionIdNotSet)?
.to_string();
Ok(ServiceBusSessionReceiver { inner, session_id })
}
}
impl<RP> ServiceBusClient<RP>
where
RP: ServiceBusRetryPolicyExt + 'static,
{
pub async fn create_rule_manager(
&mut self,
topic_name: impl AsRef<str>,
subscription_name: impl AsRef<str>,
) -> Result<ServiceBusRuleManager, azure_core::Error> {
let subscription_path = entity_name_formatter::format_subscription_path(
topic_name.as_ref(),
subscription_name.as_ref(),
);
let identifier = diagnostics::utilities::generate_identifier(&subscription_path);
let retry_options = self.connection.retry_options().clone();
let inner = self
.connection
.create_transport_rule_manager(subscription_path, identifier, retry_options)
.await?;
Ok(ServiceBusRuleManager { inner })
}
}