Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
89 changes: 88 additions & 1 deletion crates/rmcp/src/transport/common/client_side_sse.rs
Original file line number Diff line number Diff line change
Expand Up @@ -205,17 +205,24 @@ impl Default for FixedInterval {
pub struct ExponentialBackoff {
pub max_times: Option<usize>,
pub base_duration: Duration,
/// Upper bound on a single reconnect delay. The unbounded doubling policy can otherwise
/// produce delays of decades (once the multiplier saturates), which would pin the stream in
/// `tokio::time::sleep` forever — neither reconnecting nor terminating. Capping keeps the
/// backoff monotonic and panic-free while guaranteeing the client actually retries.
Comment on lines +208 to +211

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

SseAutoReconnectStream checks the retry policy before each reconnect but max_delay only limits the policy's own output.

Suggested change
/// Upper bound on a single reconnect delay. The unbounded doubling policy can otherwise
/// produce delays of decades (once the multiplier saturates), which would pin the stream in
/// `tokio::time::sleep` forever — neither reconnecting nor terminating. Capping keeps the
/// backoff monotonic and panic-free while guaranteeing the client actually retries.
/// Ceiling for the computed delay; `None` lets it grow until the multiplier saturates.
/// A server-sent `retry:` field can still push the effective wait above this.
pub max_delay: Option<Duration>,

pub max_delay: Option<Duration>,
}

impl ExponentialBackoff {
pub const DEFAULT_DURATION: Duration = Duration::from_millis(1000);
pub const DEFAULT_MAX_DELAY: Duration = Duration::from_secs(30);
}

impl Default for ExponentialBackoff {
fn default() -> Self {
Self {
max_times: None,
base_duration: Self::DEFAULT_DURATION,
max_delay: Some(Self::DEFAULT_MAX_DELAY),

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Because the struct is #[non_exhaustive], downstream callers can only create one with default() and then mutate its fields. That means anyone who sets max_times will inherit the new 30-second limit.

Do you see this shorter window as a bug fix that could go into a patch release, or would leaving max_delay set to None be the safer default?

}
}
}
Expand All @@ -227,7 +234,16 @@ impl SseRetryPolicy for ExponentialBackoff {
{
return None;
}
Some(self.base_duration * (2u32.pow(current_times as u32)))
// `current_times` is unbounded when `max_times` is unset, so the exponent can reach
// the bit width. Saturate the multiplier at `u32::MAX` and use saturating multiplication
// for the base duration so the delay stays monotonic and panic-free instead of an
// overflow panic (debug) or a wrapped-to-zero backoff (release).
let multiplier = 2u32.saturating_pow(current_times as u32);
let delay = self.base_duration.saturating_mul(multiplier);
Some(match self.max_delay {
Some(max_delay) => delay.min(max_delay),
None => delay,
})
}
}

Expand Down Expand Up @@ -775,4 +791,75 @@ mod tests {
assert!(stream.next().await.is_none());
assert_eq!(attempts.load(Ordering::Relaxed), 0);
}

#[test]
fn exponential_backoff_saturates_at_high_retry_counts() {
// With `max_times` unset, `current_times` can reach the bit width. The old
// `2u32.pow(current_times)` panicked in debug builds and wrapped in release;
// the saturating implementation must return a monotonic, non-zero delay instead.
let policy = ExponentialBackoff {
max_times: None,
base_duration: Duration::from_millis(1),
max_delay: None,
};
let mut previous = Duration::ZERO;
for current_times in [31usize, 32, 63, 64, 100] {
let delay = policy
.retry(current_times)
.expect("unbounded policy never gives up");
assert!(
!delay.is_zero(),
"delay must stay non-zero at {current_times}"
);
assert!(
delay >= previous,
"delay must stay monotonic at {current_times}"
);
previous = delay;
}
}

#[test]
fn exponential_backoff_caps_delay_at_max_delay() {
// The default cap keeps the unbounded doubling policy from producing decades-long
// sleeps once the multiplier saturates. The delay must grow monotonically, stop at
// the configured ceiling, and never exceed it.
let policy = ExponentialBackoff {
max_times: None,
base_duration: Duration::from_secs(1),
max_delay: Some(Duration::from_secs(30)),
};
let mut previous = Duration::ZERO;
for current_times in [0usize, 1, 2, 3, 4, 5, 10, 32, 64, 100] {
let delay = policy
.retry(current_times)
.expect("unbounded policy never gives up");
assert!(
delay >= previous,
"delay must stay monotonic at {current_times}"
);
assert!(
delay <= Duration::from_secs(30),
"delay must respect max_delay at {current_times}"
);
previous = delay;
}
// Beyond the ceiling the delay stays pinned at max_delay.
assert_eq!(
policy.retry(100).expect("never gives up"),
Duration::from_secs(30)
);
}

#[test]
fn exponential_backoff_respects_max_times() {
let policy = ExponentialBackoff {
max_times: Some(3),
base_duration: Duration::from_millis(1),
max_delay: None,
};
assert!(policy.retry(0).is_some());
assert!(policy.retry(2).is_some());
assert!(policy.retry(3).is_none());
}
}
2 changes: 2 additions & 0 deletions crates/rmcp/src/transport/streamable_http_client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2293,6 +2293,7 @@ mod tests {
Arc::new(ExponentialBackoff {
max_times: Some(1),
base_duration: Duration::ZERO,
max_delay: None,
}),
);
let mut stream = std::pin::pin!(stream);
Expand Down Expand Up @@ -2397,6 +2398,7 @@ mod tests {
Arc::new(ExponentialBackoff {
max_times: Some(1),
base_duration: Duration::ZERO,
max_delay: None,
}),
);

Expand Down
Loading