use core::fmt;
use core::time::Duration;
use crate::effect::{ActionDigest, ActionId, EffectKey};
use crate::error::RetryClass;
use crate::journal::EffectEvidence;
use crate::spec::RetryPolicy;
pub const MAX_SCOPE_BYTES: usize = 256;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub struct ScopeError {
len: usize,
}
impl ScopeError {
#[must_use]
pub const fn observed_len(self) -> usize {
self.len
}
}
impl fmt::Display for ScopeError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(
f,
"a deduplication scope must be at most {MAX_SCOPE_BYTES} bytes, got {}",
self.len
)
}
}
impl std::error::Error for ScopeError {}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
pub struct DedupScope(String);
impl DedupScope {
pub fn new(text: impl Into<String>) -> Result<Self, ScopeError> {
let text = text.into();
if text.len() > MAX_SCOPE_BYTES {
let refusal = Err(ScopeError { len: text.len() });
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "new: returning an error to the caller");
return refusal;
}
Ok(Self(text))
}
#[must_use]
pub fn as_str(&self) -> &str {
&self.0
}
}
impl fmt::Display for DedupScope {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str(&self.0)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
#[non_exhaustive]
pub enum LateArrival {
Refused,
Applied,
Unspecified,
}
impl LateArrival {
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
Self::Refused => "refused",
Self::Applied => "applied",
Self::Unspecified => "unspecified",
}
}
}
impl fmt::Display for LateArrival {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str(self.as_str())
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct DeduplicationContract {
action: ActionId,
payload: ActionDigest,
scope: DedupScope,
retention: Duration,
late_arrival: LateArrival,
}
impl DeduplicationContract {
#[must_use]
pub const fn new(
action: ActionId,
payload: ActionDigest,
scope: DedupScope,
retention: Duration,
late_arrival: LateArrival,
) -> Self {
Self {
action,
payload,
scope,
retention,
late_arrival,
}
}
#[must_use]
pub const fn action(&self) -> ActionId {
self.action
}
#[must_use]
pub const fn payload(&self) -> ActionDigest {
self.payload
}
#[must_use]
pub const fn scope(&self) -> &DedupScope {
&self.scope
}
#[must_use]
pub const fn retention(&self) -> Duration {
self.retention
}
#[must_use]
pub const fn late_arrival(&self) -> LateArrival {
self.late_arrival
}
#[must_use]
pub fn rules_out_duplicate(
&self,
action: ActionId,
payload: ActionDigest,
elapsed: Duration,
) -> bool {
if self.action != action || self.payload != payload {
return false;
}
if elapsed < self.retention {
return true;
}
matches!(self.late_arrival, LateArrival::Refused)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum RetryRefusal {
PermanentFailure,
BudgetExhausted {
attempts_used: u32,
max_attempts: u32,
},
AuthorityNotLive,
IntentChanged,
NotProvenUndone,
}
impl fmt::Display for RetryRefusal {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match *self {
Self::PermanentFailure => f.write_str("the failure is permanent"),
Self::BudgetExhausted {
attempts_used,
max_attempts,
} => write!(
f,
"the budget is spent: {attempts_used} of {max_attempts} attempts used"
),
Self::AuthorityNotLive => f.write_str("authority for the target is not live"),
Self::IntentChanged => {
f.write_str("the logical intent changed, so this is not a retry of that attempt")
}
Self::NotProvenUndone => f.write_str(
"nothing established that the effect did not happen, and no deduplication \
contract rules out a second application",
),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub enum RetryDecision {
Safe,
UnderContract,
Refused(Vec<RetryRefusal>),
}
impl RetryDecision {
#[must_use]
pub const fn permits_retry(&self) -> bool {
!matches!(self, Self::Refused(_))
}
#[must_use]
pub fn refusals(&self) -> &[RetryRefusal] {
match *self {
Self::Refused(ref reasons) => reasons,
Self::Safe | Self::UnderContract => &[],
}
}
}
impl fmt::Display for RetryDecision {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match *self {
Self::Safe => f.write_str("retry: the effect did not happen"),
Self::UnderContract => f.write_str("retry: a deduplication contract covers it"),
Self::Refused(ref reasons) => {
f.write_str("do not retry")?;
for (index, reason) in reasons.iter().enumerate() {
f.write_str(if index == 0 { ": " } else { "; " })?;
write!(f, "{reason}")?;
}
Ok(())
}
}
}
}
#[derive(Debug, Clone)]
pub struct RetryFacts {
policy: RetryPolicy,
attempts_used: u32,
class: RetryClass,
dispatched: EffectKey,
proposed: ActionDigest,
authority_live: bool,
evidence: Option<EffectEvidence>,
contract: Option<DeduplicationContract>,
elapsed: Duration,
}
impl RetryFacts {
#[must_use]
pub const fn new(
policy: RetryPolicy,
attempts_used: u32,
class: RetryClass,
dispatched: EffectKey,
proposed: ActionDigest,
) -> Self {
Self {
policy,
attempts_used,
class,
dispatched,
proposed,
authority_live: false,
evidence: None,
contract: None,
elapsed: Duration::ZERO,
}
}
#[must_use]
pub const fn with_live_authority(mut self) -> Self {
self.authority_live = true;
self
}
#[must_use]
pub const fn with_evidence(mut self, evidence: EffectEvidence) -> Self {
self.evidence = Some(evidence);
self
}
#[must_use]
pub fn with_contract(mut self, contract: DeduplicationContract) -> Self {
self.contract = Some(contract);
self
}
#[must_use]
pub const fn with_elapsed(mut self, elapsed: Duration) -> Self {
self.elapsed = elapsed;
self
}
fn proven_undone(&self) -> bool {
matches!(self.evidence, Some(EffectEvidence::NotApplied))
}
fn covered_by_contract(&self) -> bool {
match self.contract {
Some(ref contract) => contract.rules_out_duplicate(
self.dispatched.action(),
self.dispatched.digest(),
self.elapsed,
),
None => false,
}
}
}
#[must_use]
pub fn retry_admissible(facts: &RetryFacts) -> RetryDecision {
if facts.class == RetryClass::Never {
return RetryDecision::Refused(vec![RetryRefusal::PermanentFailure]);
}
let mut reasons = Vec::new();
if facts.attempts_used >= facts.policy.max_attempts() {
reasons.push(RetryRefusal::BudgetExhausted {
attempts_used: facts.attempts_used,
max_attempts: facts.policy.max_attempts(),
});
}
if !facts.authority_live {
reasons.push(RetryRefusal::AuthorityNotLive);
}
if facts.proposed != facts.dispatched.digest() {
reasons.push(RetryRefusal::IntentChanged);
}
let covered = facts.covered_by_contract();
if !facts.proven_undone() && !covered {
reasons.push(RetryRefusal::NotProvenUndone);
}
if !reasons.is_empty() {
return RetryDecision::Refused(reasons);
}
if facts.proven_undone() {
return RetryDecision::Safe;
}
RetryDecision::UnderContract
}
#[cfg(test)]
mod tests {
use super::*;
use crate::effect::{
ActionId, AttemptId, EnvironmentEpoch, EnvironmentId, FlowRevision, RunId,
};
const RUN: &str = "0102030405060708090a0b0c0d0e0f10";
const ACTION: &str = "1112131415161718191a1b1c1d1e1f20";
const OTHER_ACTION: &str = "3132333435363738393a3b3c3d3e3f40";
const ENV: &str = "2122232425262728292a2b2c2d2e2f30";
const FLOW_HEX: &str = "000102030405060708090a0b0c0d0e0f101112131415161718191a1b1c1d1e1f";
const DIGEST_HEX: &str = "f0f1f2f3f4f5f6f7f8f9fafbfcfdfeffe0e1e2e3e4e5e6e7e8e9eaebecedeeef";
const OTHER_DIGEST_HEX: &str =
"0102030405060708090a0b0c0d0e0f10f1f2f3f4f5f6f7f8f9fafbfcfdfeff00";
type TestResult = Result<(), Box<dyn std::error::Error>>;
fn dispatched(digest_hex: &str) -> Result<EffectKey, Box<dyn std::error::Error>> {
let key_run = RunId::from_hex(RUN)?;
let key_action = ActionId::from_hex(ACTION)?;
let key_attempt = AttemptId::from_decimal("1")?;
let key_flow_revision = FlowRevision::from_tagged("blake3_256", FLOW_HEX)?;
let key_digest = ActionDigest::from_tagged("blake3_256", digest_hex)?;
let key_environment = EnvironmentId::from_hex(ENV)?;
let key_epoch = EnvironmentEpoch::from_decimal("1")?;
Ok(
crate::effect::EffectIdentity::new(key_run, key_environment, key_flow_revision).key(
key_action,
key_attempt,
key_digest,
key_epoch,
),
)
}
fn payload() -> Result<ActionDigest, Box<dyn std::error::Error>> {
Ok(ActionDigest::from_tagged("blake3_256", DIGEST_HEX)?)
}
fn other_payload() -> Result<ActionDigest, Box<dyn std::error::Error>> {
Ok(ActionDigest::from_tagged("blake3_256", OTHER_DIGEST_HEX)?)
}
fn unresolved() -> Result<RetryFacts, Box<dyn std::error::Error>> {
Ok(RetryFacts::new(
RetryPolicy::DEFAULT,
0,
RetryClass::RequiresEvidence,
dispatched(DIGEST_HEX)?,
payload()?,
)
.with_live_authority())
}
fn permissive(class: RetryClass) -> Result<RetryFacts, Box<dyn std::error::Error>> {
Ok(RetryFacts::new(
RetryPolicy::DEFAULT,
0,
class,
dispatched(DIGEST_HEX)?,
payload()?,
)
.with_live_authority()
.with_evidence(EffectEvidence::NotApplied))
}
fn contract(
retention: Duration,
late_arrival: LateArrival,
) -> Result<DeduplicationContract, Box<dyn std::error::Error>> {
Ok(DeduplicationContract::new(
ActionId::from_hex(ACTION)?,
payload()?,
DedupScope::new("tenant-7")?,
retention,
late_arrival,
))
}
#[test]
fn proven_non_application_is_a_plain_retry() -> TestResult {
let decision = retry_admissible(&permissive(RetryClass::RequiresEvidence)?);
assert_eq!(decision, RetryDecision::Safe);
assert!(decision.permits_retry());
Ok(())
}
#[test]
fn a_permanent_failure_short_circuits() -> TestResult {
let facts = permissive(RetryClass::Never)?;
let decision = retry_admissible(&facts);
assert_eq!(
decision,
RetryDecision::Refused(vec![RetryRefusal::PermanentFailure])
);
assert!(!decision.permits_retry());
Ok(())
}
#[test]
fn a_spent_budget_refuses_and_names_the_numbers() -> TestResult {
let facts = RetryFacts::new(
RetryPolicy::DEFAULT,
3,
RetryClass::RequiresEvidence,
dispatched(DIGEST_HEX)?,
payload()?,
)
.with_live_authority()
.with_evidence(EffectEvidence::NotApplied);
assert_eq!(
retry_admissible(&facts),
RetryDecision::Refused(vec![RetryRefusal::BudgetExhausted {
attempts_used: 3,
max_attempts: 3,
}])
);
Ok(())
}
#[test]
fn one_attempt_policy_exhausts_after_the_first() -> TestResult {
let facts = RetryFacts::new(
RetryPolicy::ONE_ATTEMPT,
1,
RetryClass::Safe,
dispatched(DIGEST_HEX)?,
payload()?,
)
.with_live_authority()
.with_evidence(EffectEvidence::NotApplied);
assert!(!retry_admissible(&facts).permits_retry());
Ok(())
}
#[test]
fn absent_authority_refuses() -> TestResult {
let mut facts = permissive(RetryClass::Safe)?;
facts.authority_live = false;
assert_eq!(
retry_admissible(&facts),
RetryDecision::Refused(vec![RetryRefusal::AuthorityNotLive])
);
Ok(())
}
#[test]
fn authority_defaults_to_refusing() -> TestResult {
let facts = RetryFacts::new(
RetryPolicy::DEFAULT,
0,
RetryClass::Safe,
dispatched(DIGEST_HEX)?,
payload()?,
)
.with_evidence(EffectEvidence::NotApplied);
assert_eq!(
retry_admissible(&facts),
RetryDecision::Refused(vec![RetryRefusal::AuthorityNotLive])
);
Ok(())
}
#[test]
fn a_changed_intent_is_not_a_retry() -> TestResult {
let facts = RetryFacts::new(
RetryPolicy::DEFAULT,
0,
RetryClass::Safe,
dispatched(DIGEST_HEX)?,
other_payload()?,
)
.with_live_authority()
.with_evidence(EffectEvidence::NotApplied);
assert_eq!(
retry_admissible(&facts),
RetryDecision::Refused(vec![RetryRefusal::IntentChanged])
);
Ok(())
}
#[test]
fn an_unknown_outcome_with_no_evidence_is_refused() -> TestResult {
let facts = unresolved()?;
assert_eq!(
retry_admissible(&facts),
RetryDecision::Refused(vec![RetryRefusal::NotProvenUndone])
);
Ok(())
}
#[test]
fn a_gui_click_has_no_contract_to_reach_for() -> TestResult {
let facts = unresolved()?.with_elapsed(Duration::from_secs(1));
assert_eq!(
retry_admissible(&facts),
RetryDecision::Refused(vec![RetryRefusal::NotProvenUndone])
);
Ok(())
}
#[test]
fn a_contract_inside_its_window_covers_a_resend() -> TestResult {
let facts = unresolved()?
.with_contract(contract(Duration::from_secs(600), LateArrival::Applied)?)
.with_elapsed(Duration::from_secs(30));
assert_eq!(retry_admissible(&facts), RetryDecision::UnderContract);
Ok(())
}
#[test]
fn a_contract_past_its_window_is_worthless_if_a_late_duplicate_applies() -> TestResult {
let facts = unresolved()?
.with_contract(contract(Duration::from_secs(600), LateArrival::Applied)?)
.with_elapsed(Duration::from_secs(601));
assert_eq!(
retry_admissible(&facts),
RetryDecision::Refused(vec![RetryRefusal::NotProvenUndone])
);
Ok(())
}
#[test]
fn an_unspecified_late_arrival_is_not_a_refusal() -> TestResult {
let facts = unresolved()?
.with_contract(contract(
Duration::from_secs(600),
LateArrival::Unspecified,
)?)
.with_elapsed(Duration::from_secs(601));
assert!(
!retry_admissible(&facts).permits_retry(),
"not knowing what the remote does is not a basis for sending"
);
Ok(())
}
#[test]
fn a_contract_past_its_window_still_covers_when_the_remote_refuses_late_duplicates()
-> TestResult {
let facts = unresolved()?
.with_contract(contract(Duration::from_secs(600), LateArrival::Refused)?)
.with_elapsed(Duration::from_secs(86_400));
assert_eq!(retry_admissible(&facts), RetryDecision::UnderContract);
Ok(())
}
#[test]
fn a_contract_for_another_payload_does_not_cover() -> TestResult {
let facts = unresolved()?.with_contract(DeduplicationContract::new(
ActionId::from_hex(ACTION)?,
other_payload()?,
DedupScope::new("tenant-7")?,
Duration::from_secs(600),
LateArrival::Refused,
));
assert_eq!(
retry_admissible(&facts),
RetryDecision::Refused(vec![RetryRefusal::NotProvenUndone])
);
Ok(())
}
#[test]
fn a_contract_for_another_action_does_not_cover() -> TestResult {
let facts = unresolved()?.with_contract(DeduplicationContract::new(
ActionId::from_hex(OTHER_ACTION)?,
payload()?,
DedupScope::new("tenant-7")?,
Duration::from_secs(600),
LateArrival::Refused,
));
assert!(!retry_admissible(&facts).permits_retry());
Ok(())
}
#[test]
fn evidence_is_the_stronger_half_of_the_disjunct() -> TestResult {
let facts = RetryFacts::new(
RetryPolicy::DEFAULT,
0,
RetryClass::RequiresEvidence,
dispatched(DIGEST_HEX)?,
payload()?,
)
.with_live_authority()
.with_evidence(EffectEvidence::NotApplied)
.with_contract(contract(Duration::from_secs(600), LateArrival::Applied)?);
assert_eq!(retry_admissible(&facts), RetryDecision::Safe);
Ok(())
}
#[test]
fn applied_evidence_is_not_proof_of_non_application() -> TestResult {
let facts = RetryFacts::new(
RetryPolicy::DEFAULT,
0,
RetryClass::RequiresEvidence,
dispatched(DIGEST_HEX)?,
payload()?,
)
.with_live_authority()
.with_evidence(EffectEvidence::Applied);
assert_eq!(
retry_admissible(&facts),
RetryDecision::Refused(vec![RetryRefusal::NotProvenUndone])
);
Ok(())
}
#[test]
fn every_failing_conjunct_is_reported() -> TestResult {
let facts = RetryFacts::new(
RetryPolicy::DEFAULT,
3,
RetryClass::RequiresEvidence,
dispatched(DIGEST_HEX)?,
other_payload()?,
);
let decision = retry_admissible(&facts);
assert_eq!(
decision,
RetryDecision::Refused(vec![
RetryRefusal::BudgetExhausted {
attempts_used: 3,
max_attempts: 3,
},
RetryRefusal::AuthorityNotLive,
RetryRefusal::IntentChanged,
RetryRefusal::NotProvenUndone,
])
);
assert_eq!(decision.refusals().len(), 4);
Ok(())
}
#[test]
fn a_scope_label_is_bounded_rather_than_truncated() -> TestResult {
assert_eq!(DedupScope::new("tenant-7")?.as_str(), "tenant-7");
let long = "x".repeat(MAX_SCOPE_BYTES + 1);
match DedupScope::new(long) {
Err(error) => assert_eq!(error.observed_len(), MAX_SCOPE_BYTES + 1),
Ok(scope) => {
return Err(format!("an over-long scope must be refused, got {scope}").into());
}
}
let longest = "x".repeat(MAX_SCOPE_BYTES);
assert_eq!(DedupScope::new(longest)?.as_str().len(), MAX_SCOPE_BYTES);
Ok(())
}
#[test]
fn a_permitted_decision_reports_no_refusals() -> TestResult {
let decision = retry_admissible(&permissive(RetryClass::Safe)?);
assert!(decision.refusals().is_empty());
Ok(())
}
#[test]
fn the_refusal_renders_the_numbers_a_caller_needs() -> TestResult {
let rendered = RetryDecision::Refused(vec![
RetryRefusal::BudgetExhausted {
attempts_used: 3,
max_attempts: 3,
},
RetryRefusal::AuthorityNotLive,
])
.to_string();
assert!(rendered.contains("3 of 3 attempts used"), "{rendered}");
assert!(rendered.contains("authority"), "{rendered}");
Ok(())
}
}