use std::sync::Mutex;
use tokio::time::{self, Duration};
use tracing::{debug, error};
pub(super) enum DeliveryAttempt {
Success,
Retryable(String),
Permanent {
class: String,
reason: String,
},
}
const PERMANENT_FAILURE_REPORT_INTERVAL: Duration = Duration::from_secs(900);
#[derive(Debug, Default)]
pub(super) struct PermanentFailureReporting {
reported_class: Option<String>,
reported_at: Option<time::Instant>,
suppressed: u64,
}
impl PermanentFailureReporting {
pub(super) fn note(&mut self, class: &str, now: time::Instant) -> Option<u64> {
let due = match (self.reported_class.as_deref(), self.reported_at) {
(Some(reported_class), Some(reported_at)) => {
reported_class != class
|| now.saturating_duration_since(reported_at)
>= PERMANENT_FAILURE_REPORT_INTERVAL
}
_ => true,
};
if due {
let suppressed = self.suppressed;
self.reported_class = Some(class.to_string());
self.reported_at = Some(now);
self.suppressed = 0;
Some(suppressed)
} else {
self.suppressed += 1;
None
}
}
}
pub(super) fn report_permanent_failure(
reporting: &Mutex<PermanentFailureReporting>,
class: &str,
reason: &str,
) {
let decision = reporting
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.note(class, time::Instant::now());
match decision {
Some(suppressed) => error!(
reason = %reason,
failure_class = %class,
suppressed_since_last_report = suppressed,
report_interval_seconds = PERMANENT_FAILURE_REPORT_INTERVAL.as_secs(),
"Braintrust rejected batch; export stays off until the deployment is reconfigured"
),
None => debug!(
reason = %reason,
failure_class = %class,
"Braintrust batch rejected again"
),
}
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn a_standing_rejection_is_reported_again_after_the_interval() {
let mut reporting = PermanentFailureReporting::default();
let start = time::Instant::now();
assert_eq!(
reporting.note("HTTP 403", start),
Some(0),
"the first rejection is always reported"
);
assert_eq!(
reporting.note("HTTP 403", start + Duration::from_secs(60)),
None,
"still inside the interval: suppressed"
);
assert_eq!(
reporting.note("HTTP 403", start + Duration::from_secs(300)),
None,
"still inside the interval: suppressed"
);
assert_eq!(
reporting.note("HTTP 403", start + Duration::from_secs(1000)),
Some(2),
"past the interval it reports again, naming the two it swallowed"
);
assert_eq!(
reporting.note("HTTP 403", start + Duration::from_secs(1001)),
None,
"and the interval restarts from the report, not from the first failure"
);
}
#[tokio::test]
async fn a_different_rejection_is_not_masked_by_an_earlier_one() {
let mut reporting = PermanentFailureReporting::default();
let start = time::Instant::now();
assert_eq!(reporting.note("HTTP 403", start), Some(0));
assert_eq!(
reporting.note("HTTP 401", start + Duration::from_secs(1)),
Some(0),
"a changed failure class is reported immediately, interval or not"
);
assert_eq!(
reporting.note("HTTP 401", start + Duration::from_secs(2)),
None,
"and then it is the one being suppressed"
);
assert_eq!(
reporting.note("HTTP 403", start + Duration::from_secs(3)),
Some(1),
"switching back reports too, carrying the 401 it suppressed"
);
}
}