use std::collections::HashMap;
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
use affinidi_messaging_core::{ConnState, MessagingError, QueueFullGate};
use affinidi_tdk::messaging::ATM;
use affinidi_tdk::messaging::ReceiveHealth;
use affinidi_tdk::messaging::errors::ATMError;
use affinidi_tdk::messaging::profiles::ATMProfile;
use affinidi_tdk::messaging::protocols::message_pickup::{
MessagePickupStatusReply, UnprocessableMessage,
};
use serde::Serialize;
use tokio::sync::{broadcast, watch};
use tracing::{debug, info, warn};
pub const CHECK_INTERVAL: Duration = Duration::from_secs(30);
pub const STALE_AFTER_SECS: u64 = 30;
pub const UNANSWERED_ALARM: u32 = 3;
pub const MAX_REDELIVERY_BACKOFF: Duration = Duration::from_secs(30 * 60);
const SAME_OLDEST_TOLERANCE_SECS: u64 = 2;
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[serde(rename_all = "camelCase")]
#[non_exhaustive]
pub enum InboxCollection {
#[default]
Unknown,
Collecting,
Backlogged,
NotDelivering,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[serde(rename_all = "camelCase")]
#[non_exhaustive]
pub struct InboxHealth {
pub state: InboxCollection,
pub waiting: Option<u32>,
pub longest_waited_secs: Option<u64>,
pub unanswered_status_requests: u32,
pub unproductive_redeliveries: u32,
pub last_answered_at: Option<u64>,
pub unprocessable_seen: u64,
pub reconnects_requested: u64,
pub last_reconnect_at: Option<u64>,
pub next_reconnect_not_before: Option<u64>,
pub receive: Option<ReceiveLegHealth>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[serde(rename_all = "camelCase")]
#[non_exhaustive]
pub struct ReceiveLegHealth {
pub last_data_frame_at: Option<u64>,
pub held_frames: u32,
pub consumer_stalled_since: Option<u64>,
pub probe_outstanding_since: Option<u64>,
pub probe_reconnects: u64,
pub unprocessable_deleted: u64,
pub unprocessable_retained: u32,
}
impl From<&ReceiveHealth> for ReceiveLegHealth {
fn from(h: &ReceiveHealth) -> Self {
Self {
last_data_frame_at: h.last_data_frame_at,
held_frames: h.held_frames,
consumer_stalled_since: h.consumer_stalled_since,
probe_outstanding_since: h.probe_outstanding_since,
probe_reconnects: h.probe_reconnects,
unprocessable_deleted: h.unprocessable_deleted,
unprocessable_retained: h.unprocessable_retained,
}
}
}
#[derive(Clone, Default)]
pub struct InboxWatch {
inner: Arc<Mutex<InboxHealth>>,
receive: Arc<Mutex<Option<watch::Receiver<ReceiveHealth>>>>,
}
impl InboxWatch {
pub fn new() -> Self {
Self::default()
}
pub fn snapshot(&self) -> InboxHealth {
let mut health = self.inner.lock().expect("inbox health mutex").clone();
health.receive = self
.receive
.lock()
.expect("receive health mutex")
.as_ref()
.map(|rx| ReceiveLegHealth::from(&*rx.borrow()));
health
}
pub fn attach_receive_health(&self, rx: watch::Receiver<ReceiveHealth>) {
*self.receive.lock().expect("receive health mutex") = Some(rx);
}
fn update(&self, f: impl FnOnce(&mut InboxHealth)) {
f(&mut self.inner.lock().expect("inbox health mutex"));
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct InboxStatus {
pub waiting: u32,
pub longest_waited_secs: Option<u64>,
pub oldest_received_time: Option<u64>,
}
impl From<&MessagePickupStatusReply> for InboxStatus {
fn from(s: &MessagePickupStatusReply) -> Self {
Self {
waiting: s.message_count,
longest_waited_secs: s.longest_waited_seconds,
oldest_received_time: s.oldest_received_time,
}
}
}
impl InboxStatus {
pub fn needs_catch_up(&self) -> bool {
self.waiting > 0
&& self
.longest_waited_secs
.is_none_or(|waited| waited >= STALE_AFTER_SECS)
}
fn oldest_arrival(&self, now: u64) -> Option<u64> {
self.oldest_received_time.or_else(|| {
self.longest_waited_secs
.map(|waited| now.saturating_sub(waited))
})
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Probe {
Answered(InboxStatus),
Refused,
Unanswered,
NotSent,
}
pub fn classify_status_error(err: &ATMError) -> Probe {
match err {
ATMError::MsgSendError(msg) if msg.contains("No response") => Probe::Unanswered,
ATMError::ProblemReport(..) | ATMError::MediatorError(..) => Probe::Refused,
_ => Probe::NotSent,
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Step {
Nothing,
Redeliver,
NotDelivering,
}
#[derive(Debug, Default)]
pub struct CatchUp {
unanswered: u32,
unproductive: u32,
asked_for_oldest: Option<u64>,
redeliver_not_before: u64,
stuck_reported: bool,
not_delivering_reported: bool,
}
pub fn redelivery_backoff(unproductive: u32) -> Duration {
if unproductive == 0 {
return Duration::ZERO;
}
let factor = 1u32.checked_shl(unproductive - 1).unwrap_or(u32::MAX);
CHECK_INTERVAL
.saturating_mul(factor)
.min(MAX_REDELIVERY_BACKOFF)
}
impl CatchUp {
pub fn step(&mut self, probe: Probe, now: u64) -> Step {
match probe {
Probe::NotSent => Step::Nothing,
Probe::Unanswered => {
self.unanswered = self.unanswered.saturating_add(1);
if self.unanswered >= UNANSWERED_ALARM {
Step::NotDelivering
} else {
Step::Nothing
}
}
Probe::Refused => {
self.unanswered = 0;
self.not_delivering_reported = false;
Step::Nothing
}
Probe::Answered(status) => {
self.unanswered = 0;
self.not_delivering_reported = false;
if !status.needs_catch_up() {
self.unproductive = 0;
self.asked_for_oldest = None;
self.redeliver_not_before = 0;
self.stuck_reported = false;
return Step::Nothing;
}
let oldest = status.oldest_arrival(now);
let same_as_last_ask = match (self.asked_for_oldest, oldest) {
(Some(asked), Some(now_oldest)) => {
asked.abs_diff(now_oldest) <= SAME_OLDEST_TOLERANCE_SECS
}
_ => false,
};
if now < self.redeliver_not_before {
return Step::Nothing;
}
if same_as_last_ask {
self.unproductive = self.unproductive.saturating_add(1);
} else {
self.unproductive = 0;
self.stuck_reported = false;
}
self.asked_for_oldest = oldest;
self.redeliver_not_before =
now.saturating_add(redelivery_backoff(self.unproductive).as_secs());
Step::Redeliver
}
}
}
pub fn unanswered(&self) -> u32 {
self.unanswered
}
pub fn unproductive(&self) -> u32 {
self.unproductive
}
pub fn restart_receive_leg(&mut self) {
self.unanswered = 0;
self.not_delivering_reported = false;
}
fn take_stuck_report(&mut self) -> bool {
if self.unproductive >= 2 && !self.stuck_reported {
self.stuck_reported = true;
return true;
}
false
}
fn take_not_delivering_report(&mut self) -> bool {
if self.unanswered >= UNANSWERED_ALARM && !self.not_delivering_reported {
self.not_delivering_reported = true;
return true;
}
false
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum WatchExit {
NotDelivering,
}
pub const RECONNECT_MIN_INTERVAL: Duration = Duration::from_secs(5 * 60);
pub const RECONNECT_MAX_INTERVAL: Duration = Duration::from_secs(60 * 60);
pub const RECONNECT_REQUEST_TIMEOUT: Duration = Duration::from_secs(10);
const RECONNECT_JITTER: f64 = 0.1;
pub fn reconnect_backoff(unproductive: u32) -> Duration {
let factor = 1u32
.checked_shl(unproductive.saturating_sub(1))
.unwrap_or(u32::MAX);
RECONNECT_MIN_INTERVAL
.saturating_mul(factor)
.min(RECONNECT_MAX_INTERVAL)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Escalation {
Never,
EndSession,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Remedy {
Reconnect,
EndSession,
Wait,
}
#[derive(Debug, Default)]
struct GovernorState {
unproductive: u32,
last_at: Option<Instant>,
not_before: Option<Instant>,
last_remedy: Option<Remedy>,
requested: u64,
}
#[derive(Debug)]
pub struct ReconnectGovernor {
escalation: Escalation,
state: Mutex<GovernorState>,
}
impl ReconnectGovernor {
pub fn new(escalation: Escalation) -> Self {
Self {
escalation,
state: Mutex::default(),
}
}
pub fn escalation(&self) -> Escalation {
self.escalation
}
pub fn decide(&self, now: Instant, jitter: f64) -> Remedy {
let mut s = self.state.lock().expect("reconnect governor mutex");
if s.not_before.is_some_and(|not_before| now < not_before) {
return Remedy::Wait;
}
let remedy = match (self.escalation, s.last_remedy) {
(Escalation::EndSession, Some(Remedy::Reconnect)) => Remedy::EndSession,
_ => Remedy::Reconnect,
};
s.unproductive = s.unproductive.saturating_add(1);
s.requested = s.requested.saturating_add(1);
s.last_at = Some(now);
s.last_remedy = Some(remedy);
let wait = reconnect_backoff(s.unproductive);
let spread = wait.mul_f64(RECONNECT_JITTER * jitter.clamp(0.0, 1.0));
s.not_before = Some(now + wait + spread);
remedy
}
pub fn recovered(&self) {
let mut s = self.state.lock().expect("reconnect governor mutex");
if s.unproductive == 0 {
return;
}
s.unproductive = 0;
s.last_remedy = None;
s.not_before = s.last_at.map(|last| last + RECONNECT_MIN_INTERVAL);
}
pub fn requested(&self) -> u64 {
self.state
.lock()
.expect("reconnect governor mutex")
.requested
}
pub fn not_before(&self) -> Option<Instant> {
self.state
.lock()
.expect("reconnect governor mutex")
.not_before
}
}
fn unix_at(instant: Instant, now: Instant, unix_now: u64) -> u64 {
if instant >= now {
unix_now.saturating_add(instant.duration_since(now).as_secs())
} else {
unix_now.saturating_sub(now.duration_since(instant).as_secs())
}
}
fn unix_now() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0)
}
pub async fn run_inbox_watch(
atm: Arc<ATM>,
profile: Arc<ATMProfile>,
conn: Option<watch::Receiver<ConnState>>,
health: InboxWatch,
reconnects: Arc<ReconnectGovernor>,
) -> WatchExit {
if let Some(rx) = profile.receive_health().await {
health.attach_receive_health(rx);
}
let mut tick = tokio::time::interval(CHECK_INTERVAL);
tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
tick.tick().await;
let mut catch_up = CatchUp::default();
let alias = profile.inner.alias.clone();
loop {
tick.tick().await;
if conn
.as_ref()
.is_some_and(|c| *c.borrow() != ConnState::Connected)
{
health.update(|h| h.state = InboxCollection::Unknown);
continue;
}
let probe = match atm
.message_pickup()
.send_status_request(&profile, true, Some(Duration::from_secs(10)))
.await
{
Ok(Some(status)) => Probe::Answered(InboxStatus::from(&status)),
Ok(None) => Probe::Refused,
Err(e) => {
let probe = classify_status_error(&e);
debug!(profile = %alias, error = %e, ?probe, "inbox watch: status request failed");
probe
}
};
let was_not_delivering = catch_up.unanswered() >= UNANSWERED_ALARM;
let now = unix_now();
let step = catch_up.step(probe, now);
if matches!(probe, Probe::Answered(_) | Probe::Refused) {
reconnects.recovered();
}
health.update(|h| {
h.unanswered_status_requests = catch_up.unanswered();
h.unproductive_redeliveries = catch_up.unproductive();
match probe {
Probe::Answered(status) => {
h.waiting = Some(status.waiting);
h.longest_waited_secs = status.longest_waited_secs;
h.last_answered_at = Some(now);
h.state = if status.needs_catch_up() {
InboxCollection::Backlogged
} else {
InboxCollection::Collecting
};
}
Probe::Refused => {
h.last_answered_at = Some(now);
if h.state == InboxCollection::NotDelivering {
h.state = InboxCollection::Unknown;
}
}
Probe::Unanswered if step == Step::NotDelivering => {
h.state = InboxCollection::NotDelivering;
}
_ => {}
}
});
if was_not_delivering && catch_up.unanswered() == 0 {
info!(
profile = %alias,
"the mediator is answering again — inbound messages are reaching this node"
);
}
match step {
Step::Nothing => {}
Step::NotDelivering => {
if catch_up.take_not_delivering_report() {
warn!(
profile = %alias,
unanswered = catch_up.unanswered(),
interval_secs = CHECK_INTERVAL.as_secs(),
"the mediator socket reports connected, but status requests go \
unanswered: nothing sent to this node is reaching it, and its peers \
will be refused `limits.queue.peer` once their queue to it fills"
);
}
let at = Instant::now();
let remedy = reconnects.decide(at, rand::random::<f64>());
let not_before = reconnects
.not_before()
.map(|instant| unix_at(instant, at, now));
health.update(|h| {
h.next_reconnect_not_before = not_before;
if remedy != Remedy::Wait {
h.reconnects_requested = h.reconnects_requested.saturating_add(1);
h.last_reconnect_at = Some(now);
}
});
match remedy {
Remedy::Wait => {}
Remedy::EndSession => {
warn!(
profile = %alias,
"a reconnect did not restore delivery; ending the messaging \
session so it is rebuilt"
);
return WatchExit::NotDelivering;
}
Remedy::Reconnect => {
catch_up.restart_receive_leg();
let failure = match tokio::time::timeout(
RECONNECT_REQUEST_TIMEOUT,
profile.reconnect_websocket(),
)
.await
{
Ok(Ok(())) => {
warn!(
profile = %alias,
next_not_before = ?not_before,
"reconnecting the mediator socket so it re-registers for \
live delivery and the mediator redelivers the inbox"
);
None
}
Ok(Err(e)) => Some(format!(
"the messaging SDK has no websocket transport to reconnect \
for this profile: {e}"
)),
Err(_) => Some(format!(
"the websocket transport did not take a reconnect request \
within {}s — its task is not running its command loop",
RECONNECT_REQUEST_TIMEOUT.as_secs()
)),
};
if let Some(failure) = failure {
if reconnects.escalation() == Escalation::EndSession {
warn!(
profile = %alias,
"could not reconnect the mediator socket ({failure}); \
ending the messaging session so it is rebuilt"
);
return WatchExit::NotDelivering;
}
warn!(
profile = %alias,
next_not_before = ?not_before,
"could not reconnect the mediator socket ({failure}); the \
receive leg stays down until the transport reconnects itself"
);
}
}
}
}
Step::Redeliver => {
if let Probe::Answered(status) = probe {
if catch_up.take_stuck_report() {
warn!(
profile = %alias,
waiting = status.waiting,
longest_waited_secs = ?status.longest_waited_secs,
redeliveries = catch_up.unproductive(),
next_in_secs = redelivery_backoff(catch_up.unproductive()).as_secs(),
"messages are waiting in the mediator inbox that redelivery does not \
collect — most likely frames this node cannot unpack and the \
messaging SDK is keeping (see its 'could not unpack' warnings, \
which name their senders and why each is kept); backing off \
redelivery requests"
);
} else {
info!(
profile = %alias,
waiting = status.waiting,
longest_waited_secs = ?status.longest_waited_secs,
"mediator inbox holds messages live delivery did not bring; asking for \
redelivery"
);
}
}
if let Err(e) = atm
.message_pickup()
.toggle_live_delivery(&profile, true)
.await
{
warn!(
profile = %alias,
error = %e,
"inbox watch: could not ask the mediator to redeliver; retrying next tick"
);
}
}
}
}
}
pub async fn run_unprocessable_report(
mut rx: broadcast::Receiver<UnprocessableMessage>,
health: InboxWatch,
) {
loop {
match rx.recv().await {
Ok(frame) => {
health.update(|h| h.unprocessable_seen = h.unprocessable_seen.saturating_add(1));
debug!(
frame = frame.attachment_id.as_deref().unwrap_or("<no id>"),
reason = %frame.reason,
"the messaging SDK could not unpack an inbound message"
);
}
Err(broadcast::error::RecvError::Lagged(skipped)) => {
health.update(|h| {
h.unprocessable_seen = h.unprocessable_seen.saturating_add(skipped)
});
debug!(
skipped,
"unprocessable-frame reports lagged; counted, not described"
);
}
Err(broadcast::error::RecvError::Closed) => return,
}
}
}
pub const UNCOLLECTED_WINDOW: Duration = Duration::from_secs(5 * 60);
pub fn refused_recipient_not_collecting(err: &ATMError) -> bool {
err.http_status()
.and_then(|s| s.queue_full())
.is_some_and(|gate| gate == QueueFullGate::Peer)
}
pub fn messaging_refused_recipient_not_collecting(err: &MessagingError) -> bool {
err.queue_full() == Some(QueueFullGate::Peer)
}
#[derive(Default)]
pub struct UncollectedPeers {
marked: Mutex<HashMap<String, Instant>>,
}
impl UncollectedPeers {
pub fn new() -> Self {
Self::default()
}
pub fn record_refusal(&self, did: &str, now: Instant) -> bool {
let mut marked = self.marked.lock().expect("uncollected peers mutex");
marked.retain(|_, at| now.duration_since(*at) < UNCOLLECTED_WINDOW);
match marked.get(did) {
Some(_) => false,
None => {
marked.insert(did.to_string(), now);
true
}
}
}
pub fn record_accepted(&self, did: &str) -> bool {
self.marked
.lock()
.expect("uncollected peers mutex")
.remove(did)
.is_some()
}
pub fn is_marked(&self, did: &str, now: Instant) -> bool {
self.marked
.lock()
.expect("uncollected peers mutex")
.get(did)
.is_some_and(|at| now.duration_since(*at) < UNCOLLECTED_WINDOW)
}
pub fn observe_send(&self, did: &str, what: &str, refused_not_collecting: bool) -> bool {
if !refused_not_collecting {
return false;
}
if self.record_refusal(did, Instant::now()) {
warn!(
recipient = %did,
send = what,
window_secs = UNCOLLECTED_WINDOW.as_secs(),
"{did} is not collecting its mediator inbox; replies to it are being dropped"
);
} else {
debug!(
recipient = %did,
send = what,
"reply dropped: the recipient is still not collecting its mediator inbox"
);
}
true
}
pub fn observe_delivered(&self, did: &str) {
if self.record_accepted(did) {
info!(recipient = %did, "{did} is collecting its mediator inbox again");
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use affinidi_messaging_core::HttpStatusError;
fn status(waiting: u32, waited: Option<u64>, oldest: Option<u64>) -> Probe {
Probe::Answered(InboxStatus {
waiting,
longest_waited_secs: waited,
oldest_received_time: oldest,
})
}
#[test]
fn vti_50_only_a_message_live_delivery_missed_triggers_catch_up() {
let s = |waiting, waited| InboxStatus {
waiting,
longest_waited_secs: waited,
oldest_received_time: None,
};
assert!(!s(0, None).needs_catch_up(), "empty inbox");
assert!(!s(1, Some(2)).needs_catch_up(), "fresh push in flight");
assert!(s(1, Some(STALE_AFTER_SECS)).needs_catch_up(), "missed push");
assert!(s(1, None).needs_catch_up(), "age unreported: ask");
}
#[test]
fn vti_50_catch_up_runs_well_inside_a_short_message_expiry() {
assert!(CHECK_INTERVAL.as_secs() * 3 < 300);
}
#[test]
fn a_backlog_redelivery_does_not_clear_is_asked_for_less_and_less_often() {
let mut c = CatchUp::default();
let oldest = Some(1_000);
let mut now = 1_100;
let mut asks = Vec::new();
for _ in 0..200 {
if c.step(status(4, None, oldest), now) == Step::Redeliver {
asks.push(now);
}
now += CHECK_INTERVAL.as_secs();
}
let gaps: Vec<u64> = asks.windows(2).map(|w| w[1] - w[0]).collect();
assert!(
gaps.windows(2).all(|w| w[1] >= w[0]),
"gaps never shrink: {gaps:?}"
);
assert!(
*gaps.last().unwrap() >= MAX_REDELIVERY_BACKOFF.as_secs(),
"reaches the cap: {gaps:?}"
);
assert!(asks.len() < 20, "200 checks, {} asks", asks.len());
}
#[test]
fn a_new_backlog_is_asked_for_at_once() {
let mut c = CatchUp::default();
assert_eq!(c.step(status(1, None, Some(100)), 200), Step::Redeliver);
assert_eq!(c.step(status(1, None, Some(100)), 230), Step::Redeliver);
assert_eq!(c.unproductive(), 1);
assert_eq!(c.step(status(1, None, Some(240)), 290), Step::Redeliver);
assert_eq!(c.unproductive(), 0);
}
#[test]
fn a_drained_inbox_resets_the_backoff() {
let mut c = CatchUp::default();
for t in 0..5 {
c.step(status(1, None, Some(10)), 100 + t * 30);
}
assert!(c.unproductive() > 0);
assert_eq!(c.step(status(0, None, None), 400), Step::Nothing);
assert_eq!(c.unproductive(), 0);
assert_eq!(c.step(status(1, None, Some(390)), 430), Step::Redeliver);
}
#[test]
fn age_alone_identifies_the_same_oldest_message() {
let mut c = CatchUp::default();
assert_eq!(c.step(status(1, Some(40), None), 1_000), Step::Redeliver);
assert_eq!(c.step(status(1, Some(71), None), 1_030), Step::Redeliver);
assert_eq!(c.unproductive(), 1);
}
#[test]
fn three_unanswered_requests_while_connected_mean_not_delivering() {
let mut c = CatchUp::default();
assert_eq!(c.step(Probe::Unanswered, 0), Step::Nothing);
assert_eq!(c.step(Probe::Unanswered, 30), Step::Nothing);
assert_eq!(c.step(Probe::Unanswered, 60), Step::NotDelivering);
assert!(c.take_not_delivering_report());
assert_eq!(c.step(Probe::Unanswered, 90), Step::NotDelivering);
assert!(!c.take_not_delivering_report(), "reported once per episode");
assert_eq!(c.step(status(0, None, None), 120), Step::Nothing);
assert_eq!(c.unanswered(), 0);
}
#[test]
fn a_request_that_was_never_sent_is_not_evidence() {
let mut c = CatchUp::default();
for t in 0..10 {
assert_eq!(c.step(Probe::NotSent, t), Step::Nothing);
}
assert_eq!(c.unanswered(), 0);
}
#[test]
fn status_errors_are_classified_by_what_they_prove() {
assert_eq!(
classify_status_error(&ATMError::MsgSendError("No response from API".into())),
Probe::Unanswered
);
assert_eq!(
classify_status_error(&ATMError::TransportError(
"WebSocket message not transmitted: websocket is not connected".into()
)),
Probe::NotSent
);
assert_eq!(
classify_status_error(&ATMError::Disconnected("socket dropped".into())),
Probe::NotSent
);
assert_eq!(
classify_status_error(&ATMError::ProblemReport(
"e.p.x".into(),
"c".into(),
"".into()
)),
Probe::Refused
);
}
#[test]
fn redelivery_backoff_doubles_to_its_cap() {
assert_eq!(redelivery_backoff(0), Duration::ZERO);
assert_eq!(redelivery_backoff(1), CHECK_INTERVAL);
assert_eq!(redelivery_backoff(2), CHECK_INTERVAL * 2);
assert_eq!(redelivery_backoff(40), MAX_REDELIVERY_BACKOFF);
}
#[test]
fn reconnect_backoff_starts_at_the_floor_and_doubles_to_its_cap() {
assert_eq!(reconnect_backoff(0), RECONNECT_MIN_INTERVAL);
assert_eq!(reconnect_backoff(1), RECONNECT_MIN_INTERVAL);
assert_eq!(reconnect_backoff(2), RECONNECT_MIN_INTERVAL * 2);
assert_eq!(reconnect_backoff(3), RECONNECT_MIN_INTERVAL * 4);
assert_eq!(reconnect_backoff(40), RECONNECT_MAX_INTERVAL);
}
#[test]
fn a_dead_receive_leg_is_reconnected_at_most_once_per_interval() {
let g = ReconnectGovernor::new(Escalation::Never);
let t = Instant::now();
assert_eq!(g.decide(t, 0.0), Remedy::Reconnect);
let mut at = t;
while at < t + RECONNECT_MIN_INTERVAL {
assert_eq!(g.decide(at, 0.0), Remedy::Wait);
at += CHECK_INTERVAL;
}
assert_eq!(g.decide(t + RECONNECT_MIN_INTERVAL, 0.0), Remedy::Reconnect);
assert_eq!(g.requested(), 2);
}
#[test]
fn reconnects_that_do_not_help_back_off_to_the_cap() {
let g = ReconnectGovernor::new(Escalation::Never);
let start = Instant::now();
let mut now = start;
let mut asked = Vec::new();
while now < start + Duration::from_secs(12 * 3600) {
if g.decide(now, 1.0) == Remedy::Reconnect {
asked.push(now);
}
now += CHECK_INTERVAL;
}
let gaps: Vec<Duration> = asked.windows(2).map(|w| w[1] - w[0]).collect();
assert!(
gaps.iter().all(|g| *g >= RECONNECT_MIN_INTERVAL),
"{gaps:?}"
);
assert!(
gaps.windows(2).all(|w| w[1] >= w[0]),
"never shrinks: {gaps:?}"
);
assert!(
*gaps.last().unwrap() >= RECONNECT_MAX_INTERVAL,
"reaches the cap: {gaps:?}"
);
assert!(
*gaps.last().unwrap()
<= RECONNECT_MAX_INTERVAL + RECONNECT_MAX_INTERVAL / 10 + CHECK_INTERVAL,
"{gaps:?}"
);
assert!(asked.len() < 20, "12 h, {} reconnects", asked.len());
}
#[test]
fn a_reconnect_that_works_resets_the_backoff_but_not_the_floor() {
let g = ReconnectGovernor::new(Escalation::Never);
let t = Instant::now();
g.decide(t, 0.0);
let t2 = t + RECONNECT_MIN_INTERVAL;
g.decide(t2, 0.0);
assert_eq!(g.not_before(), Some(t2 + RECONNECT_MIN_INTERVAL * 2));
g.recovered();
assert_eq!(g.not_before(), Some(t2 + RECONNECT_MIN_INTERVAL));
assert_eq!(g.decide(t2 + Duration::from_secs(60), 0.0), Remedy::Wait);
assert_eq!(
g.decide(t2 + RECONNECT_MIN_INTERVAL, 0.0),
Remedy::Reconnect
);
}
#[test]
fn a_node_that_cannot_rebuild_its_session_never_ends_it() {
let g = ReconnectGovernor::new(Escalation::Never);
let mut now = Instant::now();
for _ in 0..50 {
assert_ne!(g.decide(now, 0.0), Remedy::EndSession);
now += RECONNECT_MAX_INTERVAL * 2;
}
}
#[test]
fn a_node_with_a_supervisor_rebuilds_its_session_when_a_reconnect_did_not_help() {
let g = ReconnectGovernor::new(Escalation::EndSession);
let t = Instant::now();
assert_eq!(g.decide(t, 0.0), Remedy::Reconnect);
let t2 = t + reconnect_backoff(1);
assert_eq!(g.decide(t2, 0.0), Remedy::EndSession);
let t3 = t2 + reconnect_backoff(2);
assert_eq!(g.decide(t3 - Duration::from_secs(1), 0.0), Remedy::Wait);
assert_eq!(g.decide(t3, 0.0), Remedy::Reconnect);
g.recovered();
assert_eq!(
g.decide(t3 + RECONNECT_MIN_INTERVAL, 0.0),
Remedy::Reconnect
);
}
#[test]
fn a_reconnected_socket_earns_its_own_strikes() {
let mut c = CatchUp::default();
for t in 0..3 {
c.step(Probe::Unanswered, t * 30);
}
assert!(c.take_not_delivering_report());
c.restart_receive_leg();
assert_eq!(c.step(Probe::Unanswered, 120), Step::Nothing);
assert_eq!(c.step(Probe::Unanswered, 150), Step::Nothing);
assert_eq!(c.step(Probe::Unanswered, 180), Step::NotDelivering);
assert!(
c.take_not_delivering_report(),
"reported again after a reconnect"
);
}
#[test]
fn instants_are_reported_as_unix_seconds() {
let now = Instant::now();
assert_eq!(unix_at(now + Duration::from_secs(300), now, 1_000), 1_300);
assert_eq!(unix_at(now, now, 1_000), 1_000);
}
#[test]
fn the_sdk_receive_view_is_reported_alongside_the_watch() {
let watch = InboxWatch::new();
assert_eq!(watch.snapshot().receive, None, "no transport, no view");
let mut sdk = ReceiveHealth::default();
sdk.last_data_frame_at = Some(1_000);
sdk.held_frames = 2;
sdk.probe_reconnects = 1;
sdk.unprocessable_deleted = 3;
sdk.unprocessable_retained = 1;
let (tx, rx) = watch::channel(sdk);
watch.attach_receive_health(rx);
let r = watch.snapshot().receive.expect("attached");
assert_eq!(r.held_frames, 2);
assert_eq!(r.unprocessable_deleted, 3);
tx.send_modify(|h| h.consumer_stalled_since = Some(1_060));
assert_eq!(
watch.snapshot().receive.unwrap().consumer_stalled_since,
Some(1_060)
);
}
#[test]
fn inbox_health_serialises_in_the_shape_diagnostics_publish() {
let watch = InboxWatch::new();
let mut sdk = ReceiveHealth::default();
sdk.probe_outstanding_since = Some(5);
let (_tx, rx) = watch::channel(sdk);
watch.attach_receive_health(rx);
let v = serde_json::to_value(watch.snapshot()).unwrap();
for key in [
"state",
"unansweredStatusRequests",
"unprocessableSeen",
"reconnectsRequested",
"lastReconnectAt",
"nextReconnectNotBefore",
"receive",
] {
assert!(v.get(key).is_some(), "missing {key}: {v}");
}
assert!(
v.get("unprocessableDeleted").is_none(),
"deletion is the SDK's, reported under `receive`: {v}"
);
let receive = &v["receive"];
for key in [
"lastDataFrameAt",
"heldFrames",
"consumerStalledSince",
"probeOutstandingSince",
"probeReconnects",
"unprocessableDeleted",
"unprocessableRetained",
] {
assert!(receive.get(key).is_some(), "missing receive.{key}: {v}");
}
assert_eq!(receive["probeOutstandingSince"], 5);
}
#[tokio::test]
async fn unprocessable_frames_are_counted_and_nothing_else() {
let (tx, rx) = broadcast::channel(4);
let watch = InboxWatch::new();
let task = tokio::spawn(run_unprocessable_report(rx, watch.clone()));
for reason in ["SecretsError: no key", "DIDComm error: bad tag"] {
tx.send(UnprocessableMessage {
attachment_id: Some("f".into()),
raw: String::new(),
reason: reason.into(),
})
.unwrap();
}
drop(tx);
task.await.unwrap();
assert_eq!(watch.snapshot().unprocessable_seen, 2);
}
fn peer_refusal() -> ATMError {
HttpStatusError::from_parts(
"send TSP message",
503,
None,
None,
r#"{"httpCode":503,"message":"{\"code\":\"e.p.limits.queue.peer\",\"comment\":\"Too many messages already waiting for this recipient\"}"}"#.to_string(),
)
.into()
}
#[test]
fn a_per_peer_refusal_means_the_recipient_is_not_collecting() {
assert!(refused_recipient_not_collecting(&peer_refusal()));
let sender_gate: ATMError = HttpStatusError::from_parts(
"send TSP message",
503,
None,
None,
r#"{"message":"{\"code\":\"e.p.limits.queue.sender\"}"}"#.to_string(),
)
.into();
assert!(!refused_recipient_not_collecting(&sender_gate));
assert!(!refused_recipient_not_collecting(
&ATMError::TransportError("connection refused".into())
));
let messaging = MessagingError::from(
peer_refusal()
.http_status()
.expect("an http status")
.clone(),
);
assert!(messaging_refused_recipient_not_collecting(&messaging));
}
#[test]
fn a_recipient_is_reported_once_per_window() {
let peers = UncollectedPeers::new();
let t = Instant::now();
assert!(peers.record_refusal("did:key:a", t), "first refusal warns");
assert!(!peers.record_refusal("did:key:a", t + Duration::from_secs(30)));
assert!(peers.record_refusal("did:key:b", t), "per recipient");
assert!(peers.is_marked("did:key:a", t + Duration::from_secs(60)));
assert!(!peers.is_marked("did:key:a", t + UNCOLLECTED_WINDOW));
assert!(peers.record_refusal("did:key:a", t + UNCOLLECTED_WINDOW + Duration::from_secs(1)));
assert!(peers.record_accepted("did:key:a"));
assert!(!peers.is_marked("did:key:a", t + UNCOLLECTED_WINDOW + Duration::from_secs(2)));
}
}