use std::time::Duration as StdDuration;
use azure_core::http::Url;
use crate::{
authorization::service_bus_token_credential::ServiceBusTokenCredential,
primitives::{
service_bus_retry_options::ServiceBusRetryOptions,
service_bus_transport_type::ServiceBusTransportType,
},
receiver::service_bus_receive_mode::ServiceBusReceiveMode,
sealed::Sealed,
};
use super::{
transport_receiver::TransportReceiver, transport_sender::TransportSender, TransportRuleManager,
TransportSessionReceiver,
};
#[cfg(docsrs)]
use crate::ServiceBusMessage;
pub(crate) trait TransportClient: Sized + Sealed {
type CreateClientError: std::error::Error + Send;
type CreateSenderError: std::error::Error + Send;
type CreateReceiverError: std::error::Error + Send;
type CreateRuleManagerError: std::error::Error + Send;
type DisposeError: std::error::Error + Send;
type Sender: TransportSender;
type Receiver: TransportReceiver;
type SessionReceiver: TransportSessionReceiver;
type RuleManager: TransportRuleManager;
async fn create_transport_client(
host: &str,
credential: ServiceBusTokenCredential,
transport_type: ServiceBusTransportType,
custom_endpoint: Option<Url>,
retry_timeout: StdDuration,
) -> Result<Self, Self::CreateClientError>;
cfg_unsecured! {
async fn create_unsecured_transport_client(
host: &str,
credential: ServiceBusTokenCredential,
transport_type: ServiceBusTransportType,
custom_endpoint: Option<Url>,
retry_timeout: StdDuration,
) -> Result<Self, Self::CreateClientError>;
}
fn transport_type(&self) -> ServiceBusTransportType;
fn is_closed(&self) -> bool;
async fn create_sender(
&mut self,
entity_path: String,
identifier: String,
retry_policy: ServiceBusRetryOptions,
) -> Result<Self::Sender, Self::CreateSenderError>;
async fn create_receiver(
&mut self,
entity_path: String,
identifier: String,
retry_options: ServiceBusRetryOptions,
receive_mode: ServiceBusReceiveMode,
prefetch_count: u32,
) -> Result<Self::Receiver, Self::CreateReceiverError>;
async fn create_session_receiver(
&mut self,
entity_path: String,
identifier: String,
retry_options: ServiceBusRetryOptions,
receive_mode: ServiceBusReceiveMode,
session_id: Option<String>,
prefetch_count: u32,
) -> Result<Self::SessionReceiver, Self::CreateReceiverError>;
async fn create_rule_manager(
&mut self,
subscription_path: String,
identifier: String,
retry_policy: ServiceBusRetryOptions,
) -> Result<Self::RuleManager, Self::CreateRuleManagerError>;
async fn close(
&mut self,
) -> Result<(), Self::DisposeError>;
async fn dispose(mut self) -> Result<(), Self::DisposeError> {
self.close().await?;
Ok(())
}
}