use fe2o3_amqp::{
connection::{self, OpenError},
link::{
DetachError, DetachThenResumeReceiverError, DetachThenResumeSenderError,
IllegalLinkStateError, LinkStateError, ReceiverAttachError, ReceiverResumeErrorKind,
RecvError, SendError, SenderAttachError, SenderResumeErrorKind,
},
session::{self, BeginError},
};
use fe2o3_amqp_management::error::{AttachError, Error as ManagementError};
use fe2o3_amqp_types::messaging::{Modified, Rejected, Released};
use timer_kit::error::Elapsed;
use crate::{
primitives::{
error::ClientDisposedError,
service_bus_retry_policy::{
should_try_recover_from_management_error, ServiceBusRetryPolicyError,
},
},
util::IntoAzureCoreError,
ServiceBusMessage,
};
#[cfg(docsrs)]
use crate::{ServiceBusPeekedMessage, ServiceBusReceivedMessage};
#[derive(Debug, thiserror::Error)]
pub(crate) enum AmqpConnectionScopeError {
#[error(transparent)]
Open(#[from] OpenError),
#[error(transparent)]
WebSocket(#[from] fe2o3_amqp_ws::Error),
#[error(transparent)]
Begin(#[from] BeginError),
#[error(transparent)]
ManagementLinkAttach(#[from] AttachError),
#[error("The connection scope is disposed")]
ScopeDisposed,
#[cfg(feature = "transaction")]
#[error(transparent)]
ControllerAttach(SenderAttachError),
}
#[derive(Debug, thiserror::Error)]
pub(crate) enum AmqpClientError {
#[error(transparent)]
UrlParseError(#[from] url::ParseError),
#[error(transparent)]
Open(#[from] OpenError),
#[error(transparent)]
WebSocket(#[from] fe2o3_amqp_ws::Error),
#[error(transparent)]
Elapsed(#[from] Elapsed),
#[error(transparent)]
Begin(#[from] BeginError),
#[error(transparent)]
ManagementLinkAttach(#[from] AttachError),
#[error(transparent)]
Dispose(#[from] DisposeError),
#[error("The client is disposed")]
ClientDisposed(#[from] ClientDisposedError),
#[cfg(feature = "transaction")]
#[error(transparent)]
ControllerAttach(SenderAttachError),
}
impl From<AmqpClientError> for azure_core::Error {
fn from(value: AmqpClientError) -> Self {
match value {
AmqpClientError::UrlParseError(error) => error.into(),
AmqpClientError::Open(error) => error.into_azure_core_error(),
AmqpClientError::WebSocket(error) => error.into_azure_core_error(),
AmqpClientError::Elapsed(error) => error.into_azure_core_error(),
AmqpClientError::Begin(error) => error.into_azure_core_error(),
AmqpClientError::ManagementLinkAttach(error) => error.into_azure_core_error(),
AmqpClientError::Dispose(error) => error.into(),
AmqpClientError::ClientDisposed(error) => error.into(),
#[cfg(feature = "transaction")]
AmqpClientError::ControllerAttach(error) => error.into_azure_core_error(),
}
}
}
impl From<AmqpConnectionScopeError> for AmqpClientError {
fn from(err: AmqpConnectionScopeError) -> Self {
match err {
AmqpConnectionScopeError::Open(err) => Self::Open(err),
AmqpConnectionScopeError::WebSocket(err) => Self::WebSocket(err),
AmqpConnectionScopeError::Begin(err) => Self::Begin(err),
AmqpConnectionScopeError::ManagementLinkAttach(err) => Self::ManagementLinkAttach(err),
AmqpConnectionScopeError::ScopeDisposed => Self::ClientDisposed(ClientDisposedError),
#[cfg(feature = "transaction")]
AmqpConnectionScopeError::ControllerAttach(err) => Self::ControllerAttach(err),
}
}
}
#[derive(Debug)]
pub struct MaxLengthExceededError {
pub(crate) message: String,
}
impl MaxLengthExceededError {
pub(crate) fn new(actual_length: usize, max_length: usize) -> Self {
Self {
message: format!(
"The actual length {} exceeds the maximum length {}",
actual_length, max_length
),
}
}
}
impl std::fmt::Display for MaxLengthExceededError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "MaxLengthExceededError: {}", self.message)
}
}
impl std::error::Error for MaxLengthExceededError {}
#[derive(Debug, thiserror::Error)]
pub enum SetMessageIdError {
#[error("Value cannot be empty")]
Empty,
#[error(transparent)]
MaxLengthExceeded(#[from] MaxLengthExceededError),
}
#[derive(Debug, thiserror::Error)]
pub enum SetPartitionKeyError {
#[error(transparent)]
MaxLengthExceeded(#[from] MaxLengthExceededError),
#[error("PartitionKey cannot be set to a different value from SessionId")]
PartitionKeyAndSessionIdAreDifferent,
}
#[derive(Debug)]
pub struct MaxAllowedTtlExceededError {}
impl std::fmt::Display for MaxAllowedTtlExceededError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(
f,
"MaxAllowedTtlExceededError: The maximum allowed TTL is u32::MAX milliseconds"
)
}
}
impl std::error::Error for MaxAllowedTtlExceededError {}
#[derive(Debug, Clone)]
pub struct RawAmqpMessageError {}
impl std::fmt::Display for RawAmqpMessageError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "RawAmqpMessageError: The message is a raw AMQP message")
}
}
impl std::error::Error for RawAmqpMessageError {}
#[derive(Debug, thiserror::Error)]
pub enum NotAcceptedError {
#[error("Rejceted: {:?}", .0)]
Rejected(Rejected),
#[error("Released: {:?}", .0)]
Released(Released),
#[error("Modified: {:?}", .0)]
Modified(Modified),
}
#[derive(Debug, thiserror::Error)]
pub(crate) enum DisposeError {
#[error(transparent)]
SessionCloseError(#[from] session::Error),
#[error(transparent)]
ConnectionCloseError(#[from] connection::Error),
}
impl From<DisposeError> for azure_core::Error {
fn from(value: DisposeError) -> Self {
match value {
DisposeError::SessionCloseError(error) => error.into_azure_core_error(),
DisposeError::ConnectionCloseError(error) => error.into_azure_core_error(),
}
}
}
#[derive(Debug, thiserror::Error)]
pub(crate) enum OpenMgmtLinkError {
#[error("Scope is disposed")]
ConnectionScopeDisposed,
#[error(transparent)]
Attach(#[from] AttachError),
#[error(transparent)]
CbsAuth(#[from] CbsAuthError),
}
#[derive(Debug, thiserror::Error)]
pub(crate) enum OpenSenderError {
#[error("The connection scope is disposed")]
ConnectionScopeDisposed,
#[error(transparent)]
ManagementLinkAttach(#[from] AttachError),
#[error(transparent)]
SenderAttach(#[from] SenderAttachError),
#[error(transparent)]
CbsAuth(#[from] CbsAuthError),
}
impl From<OpenSenderError> for azure_core::Error {
fn from(value: OpenSenderError) -> Self {
match value {
OpenSenderError::ConnectionScopeDisposed => ClientDisposedError.into(),
OpenSenderError::ManagementLinkAttach(error) => error.into_azure_core_error(),
OpenSenderError::SenderAttach(error) => error.into_azure_core_error(),
OpenSenderError::CbsAuth(error) => error.into(),
}
}
}
impl From<OpenMgmtLinkError> for OpenSenderError {
fn from(err: OpenMgmtLinkError) -> Self {
match err {
OpenMgmtLinkError::ConnectionScopeDisposed => OpenSenderError::ConnectionScopeDisposed,
OpenMgmtLinkError::Attach(err) => OpenSenderError::ManagementLinkAttach(err),
OpenMgmtLinkError::CbsAuth(err) => OpenSenderError::CbsAuth(err),
}
}
}
#[derive(Debug, thiserror::Error)]
pub(crate) enum RecoverSenderError {
#[error("The connection scope is disposed")]
ConnectionScopeDisposed,
#[error(transparent)]
Open(#[from] OpenError),
#[error(transparent)]
WebSocket(#[from] fe2o3_amqp_ws::Error),
#[error(transparent)]
Begin(#[from] BeginError),
#[error(transparent)]
ManagementLinkAttach(#[from] AttachError),
#[error(transparent)]
SenderDetach(#[from] DetachError),
#[error(transparent)]
SenderResume(#[from] SenderResumeErrorKind),
#[error(transparent)]
CbsAuth(#[from] CbsAuthError),
#[cfg(feature = "transaction")]
#[error(transparent)]
ControllerAttach(SenderAttachError),
}
impl ServiceBusRetryPolicyError for RecoverSenderError {
fn should_try_recover(&self) -> bool {
false
}
fn is_scope_disposed(&self) -> bool {
matches!(self, RecoverSenderError::ConnectionScopeDisposed)
}
}
impl From<AmqpConnectionScopeError> for RecoverSenderError {
fn from(value: AmqpConnectionScopeError) -> Self {
match value {
AmqpConnectionScopeError::Open(err) => err.into(),
AmqpConnectionScopeError::WebSocket(err) => err.into(),
AmqpConnectionScopeError::Begin(err) => err.into(),
AmqpConnectionScopeError::ManagementLinkAttach(err) => err.into(),
AmqpConnectionScopeError::ScopeDisposed => Self::ConnectionScopeDisposed,
#[cfg(feature = "transaction")]
AmqpConnectionScopeError::ControllerAttach(err) => Self::ControllerAttach(err),
}
}
}
impl From<DetachThenResumeSenderError> for RecoverSenderError {
fn from(value: DetachThenResumeSenderError) -> Self {
match value {
DetachThenResumeSenderError::Detach(err) => RecoverSenderError::SenderDetach(err),
DetachThenResumeSenderError::Resume(err) => RecoverSenderError::SenderResume(err),
}
}
}
impl From<OpenMgmtLinkError> for RecoverSenderError {
fn from(err: OpenMgmtLinkError) -> Self {
match err {
OpenMgmtLinkError::ConnectionScopeDisposed => {
RecoverSenderError::ConnectionScopeDisposed
}
OpenMgmtLinkError::Attach(err) => RecoverSenderError::ManagementLinkAttach(err),
OpenMgmtLinkError::CbsAuth(err) => RecoverSenderError::CbsAuth(err),
}
}
}
#[derive(Debug, thiserror::Error)]
pub(crate) enum OpenReceiverError {
#[error("The connection scope is disposed")]
ConnectionScopeDisposed,
#[error(transparent)]
ManagementLinkAttach(#[from] AttachError),
#[error(transparent)]
ReceiverAttach(#[from] ReceiverAttachError),
#[error(transparent)]
CbsAuth(#[from] CbsAuthError),
}
impl From<OpenReceiverError> for azure_core::Error {
fn from(value: OpenReceiverError) -> Self {
match value {
OpenReceiverError::ConnectionScopeDisposed => ClientDisposedError.into(),
OpenReceiverError::ManagementLinkAttach(error) => error.into_azure_core_error(),
OpenReceiverError::ReceiverAttach(error) => error.into_azure_core_error(),
OpenReceiverError::CbsAuth(error) => error.into(),
}
}
}
impl From<OpenMgmtLinkError> for OpenReceiverError {
fn from(err: OpenMgmtLinkError) -> Self {
match err {
OpenMgmtLinkError::ConnectionScopeDisposed => {
OpenReceiverError::ConnectionScopeDisposed
}
OpenMgmtLinkError::Attach(err) => OpenReceiverError::ManagementLinkAttach(err),
OpenMgmtLinkError::CbsAuth(err) => OpenReceiverError::CbsAuth(err),
}
}
}
#[derive(Debug, thiserror::Error)]
pub(crate) enum RecoverReceiverError {
#[error("The connection scope is disposed")]
ConnectionScopeDisposed,
#[error(transparent)]
Open(#[from] OpenError),
#[error(transparent)]
WebSocket(#[from] fe2o3_amqp_ws::Error),
#[error(transparent)]
Begin(#[from] BeginError),
#[error(transparent)]
ManagementLinkAttach(#[from] AttachError),
#[error(transparent)]
ReceiverDetach(#[from] DetachError),
#[error(transparent)]
ReceiverResume(#[from] ReceiverResumeErrorKind),
#[error(transparent)]
CbsAuth(#[from] CbsAuthError),
#[cfg(feature = "transaction")]
#[error(transparent)]
ControllerAttach(SenderAttachError),
}
impl From<AmqpConnectionScopeError> for RecoverReceiverError {
fn from(value: AmqpConnectionScopeError) -> Self {
match value {
AmqpConnectionScopeError::Open(err) => err.into(),
AmqpConnectionScopeError::WebSocket(err) => err.into(),
AmqpConnectionScopeError::Begin(err) => err.into(),
AmqpConnectionScopeError::ManagementLinkAttach(err) => err.into(),
AmqpConnectionScopeError::ScopeDisposed => Self::ConnectionScopeDisposed,
#[cfg(feature = "transaction")]
AmqpConnectionScopeError::ControllerAttach(err) => Self::ControllerAttach(err),
}
}
}
impl ServiceBusRetryPolicyError for RecoverReceiverError {
fn should_try_recover(&self) -> bool {
false
}
fn is_scope_disposed(&self) -> bool {
matches!(self, RecoverReceiverError::ConnectionScopeDisposed)
}
}
impl From<DetachThenResumeReceiverError> for RecoverReceiverError {
fn from(value: DetachThenResumeReceiverError) -> Self {
match value {
DetachThenResumeReceiverError::Detach(err) => RecoverReceiverError::ReceiverDetach(err),
DetachThenResumeReceiverError::Resume(err) => RecoverReceiverError::ReceiverResume(err),
}
}
}
impl From<OpenMgmtLinkError> for RecoverReceiverError {
fn from(err: OpenMgmtLinkError) -> Self {
match err {
OpenMgmtLinkError::ConnectionScopeDisposed => {
RecoverReceiverError::ConnectionScopeDisposed
}
OpenMgmtLinkError::Attach(err) => RecoverReceiverError::ManagementLinkAttach(err),
OpenMgmtLinkError::CbsAuth(err) => RecoverReceiverError::CbsAuth(err),
}
}
}
#[derive(Debug, thiserror::Error)]
pub(crate) enum OpenRuleManagerError {
#[error("The connection scope is disposed")]
ConnectionScopeDisposed,
#[error(transparent)]
ManagementLinkAttach(#[from] AttachError),
#[error(transparent)]
CbsAuth(#[from] CbsAuthError),
}
impl From<OpenRuleManagerError> for azure_core::Error {
fn from(value: OpenRuleManagerError) -> Self {
match value {
OpenRuleManagerError::ConnectionScopeDisposed => ClientDisposedError.into(),
OpenRuleManagerError::ManagementLinkAttach(error) => error.into_azure_core_error(),
OpenRuleManagerError::CbsAuth(error) => error.into(),
}
}
}
impl ServiceBusRetryPolicyError for OpenRuleManagerError {
fn should_try_recover(&self) -> bool {
false
}
fn is_scope_disposed(&self) -> bool {
matches!(self, OpenRuleManagerError::ConnectionScopeDisposed)
}
}
impl From<OpenMgmtLinkError> for OpenRuleManagerError {
fn from(err: OpenMgmtLinkError) -> Self {
match err {
OpenMgmtLinkError::ConnectionScopeDisposed => {
OpenRuleManagerError::ConnectionScopeDisposed
}
OpenMgmtLinkError::Attach(err) => OpenRuleManagerError::ManagementLinkAttach(err),
OpenMgmtLinkError::CbsAuth(err) => OpenRuleManagerError::CbsAuth(err),
}
}
}
#[derive(Debug, thiserror::Error)]
pub(crate) enum AmqpSendError {
#[error(transparent)]
Send(#[from] fe2o3_amqp::link::SendError),
#[error(transparent)]
NotAccepted(#[from] NotAcceptedError),
#[error(transparent)]
Elapsed(#[from] Elapsed),
}
impl ServiceBusRetryPolicyError for LinkStateError {
fn should_try_recover(&self) -> bool {
matches!(
self,
LinkStateError::IllegalState
| LinkStateError::IllegalSessionState
| LinkStateError::ExpectImmediateDetach
| LinkStateError::RemoteDetached
)
}
fn is_scope_disposed(&self) -> bool {
false
}
}
impl ServiceBusRetryPolicyError for DetachError {
fn should_try_recover(&self) -> bool {
matches!(
self,
DetachError::IllegalState
| DetachError::IllegalSessionState
| DetachError::RemoteDetachedWithError(_)
| DetachError::DetachedByRemote
)
}
fn is_scope_disposed(&self) -> bool {
false
}
}
impl ServiceBusRetryPolicyError for SendError {
fn should_try_recover(&self) -> bool {
match self {
SendError::LinkStateError(err) => err.should_try_recover(),
SendError::Detached(err) => err.should_try_recover(),
_ => false,
}
}
fn is_scope_disposed(&self) -> bool {
false
}
}
impl ServiceBusRetryPolicyError for AmqpSendError {
fn should_try_recover(&self) -> bool {
match self {
Self::Send(err) => err.should_try_recover(),
Self::Elapsed(_) => true,
_ => false,
}
}
fn is_scope_disposed(&self) -> bool {
false
}
}
#[derive(Debug, thiserror::Error)]
pub(crate) enum AmqpRecvError {
#[error(transparent)]
Recv(#[from] RecvError),
#[error(transparent)]
LinkState(#[from] IllegalLinkStateError),
#[error(transparent)]
Elapsed(#[from] Elapsed),
#[error("A valid lock token was not found in the message")]
LockTokenNotFound,
}
impl ServiceBusRetryPolicyError for AmqpRecvError {
fn should_try_recover(&self) -> bool {
matches!(
self,
Self::Recv(RecvError::LinkStateError(
LinkStateError::IllegalState
| LinkStateError::IllegalSessionState
| LinkStateError::ExpectImmediateDetach
| LinkStateError::RemoteDetached
)) | Self::LinkState(_)
)
}
fn is_scope_disposed(&self) -> bool {
false
}
}
#[derive(Debug, thiserror::Error)]
pub(crate) enum AmqpDispositionError {
#[error(transparent)]
IllegalState(#[from] IllegalLinkStateError),
#[error(transparent)]
RequestResponse(#[from] ManagementError),
#[error(transparent)]
Elapsed(#[from] Elapsed),
}
impl ServiceBusRetryPolicyError for AmqpDispositionError {
fn should_try_recover(&self) -> bool {
match self {
Self::IllegalState(IllegalLinkStateError::IllegalSessionState) => true,
Self::IllegalState(IllegalLinkStateError::IllegalState) => false,
Self::RequestResponse(err) => should_try_recover_from_management_error(err),
Self::Elapsed(_) => false,
}
}
fn is_scope_disposed(&self) -> bool {
false
}
}
#[derive(Debug, thiserror::Error)]
pub(crate) enum AmqpRequestResponseError {
#[error(transparent)]
RequestResponse(#[from] ManagementError),
#[error(transparent)]
Elapsed(#[from] Elapsed),
}
impl ServiceBusRetryPolicyError for AmqpRequestResponseError {
fn should_try_recover(&self) -> bool {
match self {
Self::RequestResponse(err) => should_try_recover_from_management_error(err),
Self::Elapsed(_) => false,
}
}
fn is_scope_disposed(&self) -> bool {
false
}
}
impl From<serde_amqp::Error> for AmqpRequestResponseError {
fn from(_: serde_amqp::Error) -> Self {
Self::RequestResponse(ManagementError::DecodeError(None))
}
}
#[derive(Debug, thiserror::Error)]
pub(crate) enum CbsAuthError {
#[error(transparent)]
TokenCredential(#[from] azure_core::Error),
#[error(transparent)]
Cbs(#[from] ManagementError),
}
impl From<CbsAuthError> for azure_core::Error {
fn from(value: CbsAuthError) -> Self {
match value {
CbsAuthError::TokenCredential(error) => error,
CbsAuthError::Cbs(error) => error.into_azure_core_error(),
}
}
}
#[derive(Debug, thiserror::Error)]
pub enum TryAddMessageError {
#[error("Message is too large to fit in a batch")]
BatchFull(ServiceBusMessage),
#[error("Cannot serialize message")]
Codec {
source: serde_amqp::Error,
message: ServiceBusMessage,
},
}
#[derive(Debug)]
pub struct RequestedSizeOutOfRange {}
impl std::fmt::Display for RequestedSizeOutOfRange {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "Requested size is out of range")
}
}
impl std::error::Error for RequestedSizeOutOfRange {}
#[derive(Debug, Clone, thiserror::Error)]
pub enum CorrelationFilterError {
#[error("Correlation filter must include at least one entry")]
EmptyFilter,
}
#[derive(Debug)]
pub(crate) struct AmqpCbsEventLoopStopped {}
impl std::fmt::Display for AmqpCbsEventLoopStopped {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "The CBS event loop has stopped")
}
}
impl std::error::Error for AmqpCbsEventLoopStopped {}
#[derive(Debug, thiserror::Error)]
pub enum CreateRuleError {
#[error("The correlation filter must have at least one entry")]
EmptyCorrelationFilter,
#[error(transparent)]
RequestResponse(#[from] ManagementError),
#[error(transparent)]
Elapsed(#[from] Elapsed),
#[error("Connection scope is disposed")]
ConnectionScopeDisposed,
}
impl From<CorrelationFilterError> for CreateRuleError {
fn from(err: CorrelationFilterError) -> Self {
match err {
CorrelationFilterError::EmptyFilter => Self::EmptyCorrelationFilter,
}
}
}
impl From<AmqpRequestResponseError> for CreateRuleError {
fn from(err: AmqpRequestResponseError) -> Self {
match err {
AmqpRequestResponseError::RequestResponse(err) => Self::RequestResponse(err),
AmqpRequestResponseError::Elapsed(err) => Self::Elapsed(err),
}
}
}
impl ServiceBusRetryPolicyError for CreateRuleError {
fn should_try_recover(&self) -> bool {
match self {
Self::RequestResponse(err) => should_try_recover_from_management_error(err),
Self::Elapsed(_) => false,
Self::ConnectionScopeDisposed => false,
Self::EmptyCorrelationFilter => false,
}
}
fn is_scope_disposed(&self) -> bool {
matches!(self, Self::ConnectionScopeDisposed)
}
}