use super::TransportError;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)]
pub enum FailureAction {
Nack,
DeadLetter,
Park,
LogAndAck,
Stop,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)]
pub enum FailurePolicy {
Retry,
DeadLetter,
Park,
LogAndAck,
Stop,
}
impl Default for FailurePolicy {
fn default() -> Self {
FailurePolicy::DeadLetter
}
}
impl FailurePolicy {
pub fn resolve(self, error: &TransportError) -> FailureAction {
if error.is_retryable() {
return FailureAction::Nack;
}
match self {
FailurePolicy::Retry => FailureAction::Nack,
FailurePolicy::DeadLetter => FailureAction::DeadLetter,
FailurePolicy::Park => FailureAction::Park,
FailurePolicy::LogAndAck => FailureAction::LogAndAck,
FailurePolicy::Stop => FailureAction::Stop,
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn default_policy_dead_letters_permanent_failures() {
assert_eq!(FailurePolicy::default(), FailurePolicy::DeadLetter);
let action = FailurePolicy::default().resolve(&TransportError::permanent("bad"));
assert_eq!(action, FailureAction::DeadLetter);
}
#[test]
fn retryable_errors_always_nack_regardless_of_policy() {
let retry = TransportError::retryable("transient");
for policy in [
FailurePolicy::Retry,
FailurePolicy::DeadLetter,
FailurePolicy::Park,
FailurePolicy::LogAndAck,
FailurePolicy::Stop,
] {
assert_eq!(
policy.resolve(&retry),
FailureAction::Nack,
"{policy:?} should nack a retryable error"
);
}
}
#[test]
fn permanent_failures_follow_the_configured_policy() {
let permanent = TransportError::permanent("terminal");
assert_eq!(
FailurePolicy::Retry.resolve(&permanent),
FailureAction::Nack
);
assert_eq!(
FailurePolicy::DeadLetter.resolve(&permanent),
FailureAction::DeadLetter
);
assert_eq!(FailurePolicy::Park.resolve(&permanent), FailureAction::Park);
assert_eq!(
FailurePolicy::LogAndAck.resolve(&permanent),
FailureAction::LogAndAck
);
assert_eq!(FailurePolicy::Stop.resolve(&permanent), FailureAction::Stop);
}
#[test]
fn only_log_and_ack_discards_a_failed_message() {
let permanent = TransportError::permanent("terminal");
let discards =
|policy: FailurePolicy| policy.resolve(&permanent) == FailureAction::LogAndAck;
assert!(discards(FailurePolicy::LogAndAck));
assert!(!discards(FailurePolicy::Retry));
assert!(!discards(FailurePolicy::DeadLetter));
assert!(!discards(FailurePolicy::Park));
assert!(!discards(FailurePolicy::Stop));
}
}