use std::time::Duration as StdDuration;
use fe2o3_amqp::link::SendError;
use fe2o3_amqp_management::error::Error as ManagementError;
use super::service_bus_retry_options::ServiceBusRetryOptions;
pub(crate) static SERVER_BUSY_BASE_SLEEP_TIME: StdDuration = StdDuration::from_secs(10);
pub trait ServiceBusRetryPolicyError: std::error::Error + Send + Sync + 'static {
fn should_try_recover(&self) -> bool;
fn is_scope_disposed(&self) -> bool;
}
pub(crate) fn should_try_recover_from_management_error(
error: &fe2o3_amqp_management::error::Error,
) -> bool {
use fe2o3_amqp::link::{LinkStateError, RecvError};
matches!(
error,
ManagementError::Send(SendError::LinkStateError(
LinkStateError::IllegalSessionState
)) | ManagementError::Recv(RecvError::LinkStateError(
LinkStateError::IllegalSessionState
))
)
}
pub trait ServiceBusRetryPolicy: std::fmt::Debug + Send {
fn options(&self) -> &ServiceBusRetryOptions;
fn state(&self) -> &dyn ServiceBusRetryPolicyState;
fn state_mut(&mut self) -> &mut dyn ServiceBusRetryPolicyState;
fn calculate_try_timeout(&self, attempt_count: u32) -> StdDuration;
fn calculate_retry_delay(
&self,
last_error: &dyn ServiceBusRetryPolicyError,
attempt_count: u32,
) -> Option<StdDuration>;
}
pub trait ServiceBusRetryPolicyState {
fn is_server_busy(&self) -> bool;
fn set_server_busy(&mut self, error_message: String);
fn reset_server_busy(&mut self);
fn server_busy_error_message(&self) -> Option<&str>;
}
pub trait ServiceBusRetryPolicyExt:
ServiceBusRetryPolicy + From<ServiceBusRetryOptions> + Send + Sync
{
}
impl<T> ServiceBusRetryPolicyExt for T where
T: ServiceBusRetryPolicy + From<ServiceBusRetryOptions> + Send + Sync
{
}
macro_rules! run_operation {
($policy:tt, $err_ty:ty, $try_timeout:ident, $op:expr, $recover_op:expr) => {{
let mut _failed_attempt_count = 0; let mut _should_try_recover = false; let mut _is_scope_disposed = false;
if crate::primitives::service_bus_retry_policy::ServiceBusRetryPolicyState::is_server_busy($policy.state())
&& $try_timeout
< crate::primitives::service_bus_retry_policy::SERVER_BUSY_BASE_SLEEP_TIME
{
let _sleep_result = crate::util::time::sleep($try_timeout).await;
crate::util::time::handle_delay_value(_sleep_result, &mut _is_scope_disposed);
}
let outcome = loop {
if crate::primitives::service_bus_retry_policy::ServiceBusRetryPolicyState::is_server_busy($policy.state()) {
let _sleep_result = crate::util::time::sleep(
crate::primitives::service_bus_retry_policy::SERVER_BUSY_BASE_SLEEP_TIME,
)
.await;
crate::util::time::handle_delay_value(_sleep_result, &mut _is_scope_disposed);
}
if _should_try_recover {
match crate::util::time::timeout($try_timeout, $recover_op).await {
Ok(result) => match result {
Ok(_) => {
log::info!("Recovery succeeded");
_should_try_recover = false;
}
Err(recover_error) => {
log::error!("Failed to recover {}", &recover_error);
_is_scope_disposed |= crate::primitives::service_bus_retry_policy::ServiceBusRetryPolicyError::is_scope_disposed(&recover_error);
_should_try_recover = crate::primitives::service_bus_retry_policy::ServiceBusRetryPolicyError::should_try_recover(&recover_error);
}
}
Err(elapsed) => {
let err = <$err_ty>::from(elapsed);
_is_scope_disposed |= crate::primitives::service_bus_retry_policy::ServiceBusRetryPolicyError::is_scope_disposed(&err);
}
}
}
let outcome = match crate::util::time::timeout($try_timeout, $op).await {
Ok(result) => result.map_err(<$err_ty>::from),
Err(elapsed) => Err(<$err_ty>::from(elapsed)),
};
match outcome {
Ok(outcome) => break outcome,
Err(error) => {
_failed_attempt_count += 1;
let _retry_delay = $policy.calculate_retry_delay(&error, _failed_attempt_count);
_is_scope_disposed |= crate::primitives::service_bus_retry_policy::ServiceBusRetryPolicyError::is_scope_disposed(&error);
_should_try_recover = crate::primitives::service_bus_retry_policy::ServiceBusRetryPolicyError::should_try_recover(&error);
match (_retry_delay, _is_scope_disposed) {
(Some(retry_delay), false) => {
log::error!("{}", &error);
let _sleep_result = crate::util::time::sleep(retry_delay).await;
crate::util::time::handle_delay_value(_sleep_result, &mut _is_scope_disposed);
$try_timeout = $policy.calculate_try_timeout(_failed_attempt_count);
}
_ => return Err(crate::primitives::error::RetryError::Operation(error)),
}
}
}
};
Result::<_, crate::primitives::error::RetryError<$err_ty>>::Ok(outcome)
}};
}
pub(crate) use run_operation;