use std::time::Duration;
use tokio_util::sync::CancellationToken;
use crate::client::CyclesClient;
use crate::error::Error;
use crate::models::enums::ErrorCode;
use crate::models::{ExtendRequest, IdempotencyKey, ReservationId};
const HELD_CADENCE_CAP_MS: u64 = 30_000;
const MIN_BEAT_DELAY_MS: u64 = 500;
fn held_cadence_ms(requested_ttl_ms: u64) -> u64 {
(requested_ttl_ms / 2).clamp(1, HELD_CADENCE_CAP_MS)
}
fn is_lead_clamp_grant(grant_ms: i128, requested_ttl_ms: u64, elapsed_ms: u128) -> bool {
if grant_ms <= 0 {
return true;
}
let elapsed = i128::try_from(elapsed_ms).unwrap_or(i128::MAX / 8);
10 * grant_ms < 9 * i128::from(requested_ttl_ms)
&& 4 * grant_ms <= 5 * elapsed
&& 4 * grant_ms >= 3 * elapsed
}
fn next_beat_delay_ms(last_grant_ms: u64, requested_ttl_ms: u64) -> u64 {
let hi = (requested_ttl_ms / 2).max(MIN_BEAT_DELAY_MS);
(last_grant_ms / 2).clamp(MIN_BEAT_DELAY_MS, hi)
}
fn lead_min_ms(grants_sum_ms: i128, elapsed_ms: u128) -> i128 {
grants_sum_ms - elapsed_ms as i128
}
fn should_skip(lead_min_ms: i128, last_grant_ms: Option<u64>) -> bool {
match last_grant_ms {
Some(grant) => 2 * lead_min_ms >= 3 * i128::from(grant),
None => false,
}
}
fn is_permanent_extend_failure(e: &Error) -> bool {
matches!(
e.error_code(),
Some(
ErrorCode::ReservationExpired
| ErrorCode::ReservationFinalized
| ErrorCode::MaxExtensionsExceeded
| ErrorCode::TenantClosed
| ErrorCode::NotFound
)
) || matches!(e.status(), Some(410 | 404))
}
fn advance(intended: tokio::time::Instant, delay: Duration) -> tokio::time::Instant {
(intended + delay).max(tokio::time::Instant::now())
}
pub(crate) fn ceil_ms(d: Duration) -> u64 {
u64::try_from(d.as_nanos().div_ceil(1_000_000)).unwrap_or(u64::MAX)
}
fn lead_floor_ms(remaining_ttl_ms: u64, rtt_ms: u64) -> u64 {
remaining_ttl_ms.saturating_sub(rtt_ms)
}
fn attempt_budget_ms(request_timeout_ms: Option<u64>, max_rtt_ms: u64) -> Option<u64> {
request_timeout_ms.map(|t| t.max(1_000).max(max_rtt_ms.saturating_mul(2)))
}
fn safety_margin_ms(max_rtt_ms: u64) -> u64 {
max_rtt_ms.saturating_mul(2).max(1_000)
}
fn retry_reserve_ms(request_timeout_ms: Option<u64>, max_rtt_ms: u64) -> Option<u64> {
attempt_budget_ms(request_timeout_ms, max_rtt_ms).map(|b| {
b.saturating_mul(2)
.saturating_add(safety_margin_ms(max_rtt_ms))
})
}
fn remaining_next_delay_ms(
remaining_ttl_ms: u64,
rtt_ms: u64,
request_timeout_ms: Option<u64>,
max_rtt_ms: u64,
) -> u64 {
match retry_reserve_ms(request_timeout_ms, max_rtt_ms) {
Some(reserve) => lead_floor_ms(remaining_ttl_ms, rtt_ms).saturating_sub(reserve),
None => 0,
}
}
fn recovery_window_ms(
lead_estimate_ms: u64,
request_timeout_ms: Option<u64>,
max_rtt_ms: u64,
) -> Option<i128> {
attempt_budget_ms(request_timeout_ms, max_rtt_ms).map(|budget| {
i128::from(lead_estimate_ms) - i128::from(budget) - i128::from(safety_margin_ms(max_rtt_ms))
})
}
fn recovery_delay_ms(lead_estimate_ms: u64, window_ms: i128) -> u64 {
let window = u64::try_from(window_ms).unwrap_or(0);
30_000.min(lead_estimate_ms / 4).min(window)
}
fn recovery_stalled(
window_ms: i128,
previous_window_ms: Option<i128>,
elapsed_advanced: bool,
) -> bool {
previous_window_ms.is_some_and(|prev| !elapsed_advanced && window_ms >= prev)
}
fn is_recoverable_in_field_mode(e: &Error) -> bool {
matches!(e, Error::Transport(_) | Error::Deserialization(_))
|| matches!(e.status(), Some(s) if s >= 500 || s == 429)
}
pub(crate) struct CreateLeaseSample {
pub(crate) remaining_ttl_ms: Option<u64>,
pub(crate) rtt_ms: u64,
pub(crate) received_at: tokio::time::Instant,
}
pub(crate) fn start_heartbeat(
client: CyclesClient,
reservation_id: ReservationId,
requested_ttl_ms: u64,
initial_expires_at_ms: Option<u64>,
create_sample: CreateLeaseSample,
cancel: CancellationToken,
) -> tokio::task::JoinHandle<()> {
tokio::spawn(async move {
let initial_remaining_ttl_ms = create_sample.remaining_ttl_ms;
let create_rtt_ms = create_sample.rtt_ms;
let create_received_at = create_sample.received_at;
let anchor = tokio::time::Instant::now();
let mut delay = Duration::from_millis(held_cadence_ms(requested_ttl_ms));
let request_timeout_ms: Option<u64> = Some(ceil_ms(
client.config().connect_timeout + client.config().read_timeout,
));
let mut normative = initial_remaining_ttl_ms.is_some();
let mut max_rtt_ms = create_rtt_ms;
let elapsed_after_receipt_ms =
ceil_ms(anchor.saturating_duration_since(create_received_at));
let remaining_at_start =
initial_remaining_ttl_ms.map(|r| r.saturating_sub(elapsed_after_receipt_ms));
let mut lead_sample: Option<(u64, tokio::time::Instant)> =
remaining_at_start.map(|r| (lead_floor_ms(r, create_rtt_ms), anchor));
let mut zero_delay_streak: u32 = 0;
let mut last_recovery_window: Option<i128> = None;
let mut last_recovery_failure_at: Option<tokio::time::Instant> = None;
let mut zero_window_retried = false;
let mut next_beat = match remaining_at_start {
Some(remaining) => {
let nd = remaining_next_delay_ms(
remaining,
create_rtt_ms,
request_timeout_ms,
max_rtt_ms,
);
if nd == 0 {
zero_delay_streak = 1;
}
anchor + Duration::from_millis(nd)
}
None => anchor,
};
let mut prev_expiry = initial_expires_at_ms;
let mut grants_sum_ms: i128 = 0;
let mut last_grant_ms: Option<u64> = None;
let mut last_success = anchor;
let mut warned_lead_clamp = false;
let mut pending_key: Option<IdempotencyKey> = None;
loop {
tokio::select! {
() = cancel.cancelled() => break,
() = tokio::time::sleep_until(next_beat) => {}
}
let elapsed_ms = next_beat.duration_since(anchor).as_millis();
if !normative && should_skip(lead_min_ms(grants_sum_ms, elapsed_ms), last_grant_ms) {
next_beat = advance(next_beat, delay);
continue;
}
let key = pending_key.clone().unwrap_or_else(IdempotencyKey::random);
let req = ExtendRequest {
idempotency_key: key.clone(),
extend_by_ms: requested_ttl_ms,
metadata: None,
};
let sent_at = tokio::time::Instant::now();
let result = client
.extend_reservation_strict(&reservation_id, &req)
.await;
match result {
Ok(resp) => {
let received_at = tokio::time::Instant::now();
let rtt_ms = ceil_ms(received_at.duration_since(sent_at));
max_rtt_ms = max_rtt_ms.max(rtt_ms);
pending_key = None;
last_recovery_window = None;
last_recovery_failure_at = None;
zero_window_retried = false;
let grant: i128 = match prev_expiry {
Some(prev) => i128::from(resp.expires_at_ms) - i128::from(prev),
None => i128::from(requested_ttl_ms),
};
prev_expiry = Some(resp.expires_at_ms);
let elapsed_since_success = next_beat.duration_since(last_success).as_millis();
last_success = next_beat;
let counted = u64::try_from(grant.max(0)).unwrap_or(u64::MAX);
grants_sum_ms += i128::from(counted);
last_grant_ms = Some(counted);
match resp.remaining_ttl_ms {
Some(remaining) => {
normative = true;
let floor = lead_floor_ms(remaining, rtt_ms);
lead_sample = Some((floor, received_at));
let nd = remaining_next_delay_ms(
remaining,
rtt_ms,
request_timeout_ms,
max_rtt_ms,
);
if nd == 0 {
zero_delay_streak += 1;
if zero_delay_streak >= 2 {
tracing::warn!(
reservation_id = %reservation_id,
remaining_ttl_ms = remaining,
retry_reserve_ms =
?retry_reserve_ms(request_timeout_ms, max_rtt_ms),
"reservation lease is shorter than the heartbeat's retry-safety budget (next_delay = 0 twice in a row); stopping heartbeat"
);
break;
}
next_beat = received_at;
} else {
zero_delay_streak = 0;
next_beat = received_at + Duration::from_millis(nd);
}
}
None => {
normative = false;
zero_delay_streak = 0;
if is_lead_clamp_grant(grant, requested_ttl_ms, elapsed_since_success) {
delay = Duration::from_millis(held_cadence_ms(requested_ttl_ms));
if !warned_lead_clamp {
warned_lead_clamp = true;
tracing::warn!(
reservation_id = %reservation_id,
grant_ms = %grant,
elapsed_ms = %elapsed_since_success,
requested_ttl_ms,
held_cadence_ms = held_cadence_ms(requested_ttl_ms),
"extend grants track elapsed time, not the requested lease — server appears to clamp the reservation's maximum lead; holding heartbeat cadence to avoid depleting the extension allowance"
);
}
} else {
delay = Duration::from_millis(next_beat_delay_ms(
counted,
requested_ttl_ms,
));
}
next_beat = advance(next_beat, delay);
}
}
}
Err(e) if is_permanent_extend_failure(&e) => {
tracing::warn!(
reservation_id = %reservation_id,
error = %e,
"heartbeat extend failed permanently; stopping heartbeat"
);
break;
}
Err(e) if normative => {
if !is_recoverable_in_field_mode(&e) {
tracing::warn!(
reservation_id = %reservation_id,
error = %e,
"heartbeat extend rejected (non-retryable request/authorization failure); stopping heartbeat without key rotation"
);
break;
}
let failure_at = tokio::time::Instant::now();
let lead_now = lead_sample.map_or(0, |(floor, at)| {
floor.saturating_sub(ceil_ms(failure_at.duration_since(at)))
});
let window = recovery_window_ms(lead_now, request_timeout_ms, max_rtt_ms)
.unwrap_or(i128::MIN);
if window == 0 && zero_window_retried {
tracing::warn!(
reservation_id = %reservation_id,
error = %e,
"heartbeat retry window is still zero after the one permitted immediate recovery retry; stopping heartbeat"
);
break;
}
let elapsed_advanced =
last_recovery_failure_at.is_none_or(|previous| failure_at > previous);
if recovery_stalled(window, last_recovery_window, elapsed_advanced) {
tracing::warn!(
reservation_id = %reservation_id,
error = %e,
"heartbeat recovery made no progress between consecutive failures; stopping heartbeat"
);
break;
}
if window < 0 {
tracing::warn!(
reservation_id = %reservation_id,
error = %e,
lead_estimate_ms = lead_now,
"no complete extend retry plus safety margin fits the remaining lease; stopping heartbeat (lease cannot be safely renewed)"
);
break;
}
if window == 0 {
zero_window_retried = true;
}
let retry_delay_ms = if e.status() == Some(429) {
match e.retry_after().map(ceil_ms) {
Some(ra_ms) if i128::from(ra_ms) <= window => ra_ms,
ra => {
tracing::warn!(
reservation_id = %reservation_id,
error = %e,
retry_after_ms = ?ra,
"429 Retry-After is missing, invalid, or exceeds the safe retry window; stopping heartbeat (lease cannot be safely renewed)"
);
break;
}
}
} else {
recovery_delay_ms(lead_now, window)
};
pending_key = Some(key);
last_recovery_window = Some(window);
last_recovery_failure_at = Some(failure_at);
next_beat = tokio::time::Instant::now() + Duration::from_millis(retry_delay_ms);
tracing::warn!(
reservation_id = %reservation_id,
error = %e,
retry_delay_ms,
"heartbeat extend failed; retrying with the same idempotency key within the recovery window"
);
}
Err(e) => {
pending_key = Some(key);
next_beat = advance(next_beat, delay);
tracing::warn!(
reservation_id = %reservation_id,
error = %e,
"heartbeat extend failed; retrying next beat with the same idempotency key"
);
}
}
}
})
}
#[cfg(test)]
mod tests {
use super::*;
const REQUESTED: u64 = 2_000;
#[test]
fn lead_min_starts_at_zero_and_goes_negative() {
assert_eq!(lead_min_ms(0, 0), 0);
assert_eq!(lead_min_ms(0, 750), -750);
assert_eq!(lead_min_ms(4_000, 3_000), 1_000);
assert_eq!(lead_min_ms(2_000, 10_000), -8_000);
}
#[test]
fn no_skip_before_first_grant_sample() {
assert!(!should_skip(i128::MAX / 4, None));
assert!(!should_skip(0, None));
}
#[test]
fn skip_threshold_is_1_5_times_last_grant_inclusive() {
assert!(should_skip(3_000, Some(REQUESTED)));
assert!(!should_skip(2_999, Some(REQUESTED)));
assert!(!should_skip(-1, Some(REQUESTED)));
assert!(should_skip(0, Some(0)));
assert!(!should_skip(-1, Some(0)));
}
#[test]
fn full_grant_trace_extends_three_beats_then_alternates() {
let mut grants: i128 = 0;
let mut extends = Vec::new();
for beat in 1u128..=6 {
let lead = lead_min_ms(grants, (beat - 1) * 1_000);
let skip = should_skip(lead, if beat == 1 { None } else { Some(REQUESTED) });
extends.push(!skip);
if !skip {
grants += i128::from(REQUESTED);
}
}
assert_eq!(extends, [true, true, true, false, true, false]);
}
#[test]
fn held_cadence_pins() {
assert_eq!(held_cadence_ms(2_000), 1_000);
assert_eq!(held_cadence_ms(1_000), 500);
assert_eq!(held_cadence_ms(86_400_000), 30_000);
assert_eq!(held_cadence_ms(60_000), 30_000);
assert_eq!(held_cadence_ms(1), 1);
}
#[test]
fn lead_clamp_zero_or_negative_grant_always_holds() {
assert!(is_lead_clamp_grant(0, REQUESTED, 0));
assert!(is_lead_clamp_grant(-500, REQUESTED, 1_000));
assert!(is_lead_clamp_grant(0, REQUESTED, 10_000));
}
#[test]
fn lead_clamp_full_grant_is_trusted_regardless_of_elapsed() {
assert!(!is_lead_clamp_grant(2_000, REQUESTED, 2_000));
assert!(!is_lead_clamp_grant(1_800, REQUESTED, 1_800));
assert!(is_lead_clamp_grant(1_799, REQUESTED, 1_799));
}
#[test]
fn lead_clamp_band_is_0_75_to_1_25_of_elapsed() {
const REQ: u64 = 8_000;
assert!(is_lead_clamp_grant(1_000, REQ, 1_000));
assert!(is_lead_clamp_grant(1_250, REQ, 1_000));
assert!(!is_lead_clamp_grant(1_251, REQ, 1_000));
assert!(is_lead_clamp_grant(750, REQ, 1_000));
assert!(!is_lead_clamp_grant(749, REQ, 1_000));
assert!(!is_lead_clamp_grant(1_000, REQ, 2_000));
assert!(!is_lead_clamp_grant(2_000, REQ, 1_000));
assert!(!is_lead_clamp_grant(5, REQ, 0));
}
#[test]
fn next_beat_delay_tracks_grant_within_bounds() {
assert_eq!(next_beat_delay_ms(2_000, 2_000), 1_000);
assert_eq!(next_beat_delay_ms(500, 2_000), 500);
assert_eq!(next_beat_delay_ms(0, 2_000), 500);
assert_eq!(next_beat_delay_ms(20_000, 8_000), 4_000);
assert_eq!(next_beat_delay_ms(2_000, 8_000), 1_000);
}
#[test]
fn lead_floor_subtracts_rtt_and_saturates() {
assert_eq!(lead_floor_ms(2_000, 3), 1_997);
assert_eq!(lead_floor_ms(2_000, 0), 2_000);
assert_eq!(lead_floor_ms(100, 200), 0);
}
#[test]
fn ceil_ms_rounds_up() {
assert_eq!(ceil_ms(Duration::from_millis(5)), 5);
assert_eq!(ceil_ms(Duration::from_micros(1)), 1);
assert_eq!(ceil_ms(Duration::from_micros(4_001)), 5);
assert_eq!(ceil_ms(Duration::ZERO), 0);
}
#[test]
fn attempt_budget_and_safety_margin_pins() {
assert_eq!(attempt_budget_ms(Some(500), 0), Some(1_000));
assert_eq!(attempt_budget_ms(Some(10_000), 1_500), Some(10_000));
assert_eq!(attempt_budget_ms(Some(1_000), 5_000), Some(10_000));
assert_eq!(attempt_budget_ms(None, 5_000), None);
assert_eq!(safety_margin_ms(0), 1_000);
assert_eq!(safety_margin_ms(400), 1_000);
assert_eq!(safety_margin_ms(1_500), 3_000);
}
#[test]
fn remaining_next_delay_pins() {
assert_eq!(
remaining_next_delay_ms(60_000, 0, Some(10_000), 1_500),
37_000
);
assert_eq!(remaining_next_delay_ms(60_000, 0, Some(30_000), 0), 0);
assert_eq!(remaining_next_delay_ms(60_000, 0, Some(500), 0), 57_000);
assert_eq!(remaining_next_delay_ms(4_000, 0, Some(500), 0), 1_000);
assert_eq!(remaining_next_delay_ms(3_000, 0, Some(500), 0), 0);
assert_eq!(remaining_next_delay_ms(4_000, 500, Some(500), 0), 500);
assert_eq!(remaining_next_delay_ms(3_600_000, 0, None, 0), 0);
assert_eq!(remaining_next_delay_ms(0, 0, Some(500), 0), 0);
assert_eq!(remaining_next_delay_ms(100, 200, Some(500), 0), 0);
}
#[test]
fn recovery_window_pins() {
assert_eq!(recovery_window_ms(2_985, Some(500), 0), Some(985));
assert_eq!(recovery_window_ms(2_000, Some(500), 0), Some(0));
assert_eq!(recovery_window_ms(1_997, Some(500), 0), Some(-3));
assert_eq!(recovery_window_ms(0, Some(500), 0), Some(-2_000));
assert_eq!(recovery_window_ms(10_000, Some(1_000), 1_500), Some(4_000));
assert_eq!(recovery_window_ms(u64::MAX, None, 0), None);
}
#[test]
fn recovery_delay_is_min_of_cap_quarter_lead_and_window() {
assert_eq!(recovery_delay_ms(2_985, 985), 746);
assert_eq!(recovery_delay_ms(8_000, 30_000), 2_000);
assert_eq!(recovery_delay_ms(400_000, 200_000), 30_000);
assert_eq!(recovery_delay_ms(2_000, 0), 0);
}
#[test]
fn recovery_progress_guard() {
assert!(!recovery_stalled(985, None, false));
assert!(!recovery_stalled(0, None, false));
assert!(!recovery_stalled(240, Some(985), false));
assert!(!recovery_stalled(-3, Some(240), false));
assert!(!recovery_stalled(985, Some(985), true));
assert!(recovery_stalled(985, Some(985), false));
assert!(recovery_stalled(986, Some(985), false));
assert!(recovery_stalled(0, Some(0), false));
}
#[test]
fn field_mode_recoverable_classification() {
let api = |status: u16| Error::Api {
status,
code: None,
message: "x".into(),
request_id: None,
retry_after: None,
details: None,
};
assert!(is_recoverable_in_field_mode(&api(500)));
assert!(is_recoverable_in_field_mode(&api(503)));
assert!(is_recoverable_in_field_mode(&api(429)));
assert!(is_recoverable_in_field_mode(&Error::Deserialization(
serde::de::Error::custom("ambiguous")
)));
assert!(!is_recoverable_in_field_mode(&api(400)));
assert!(!is_recoverable_in_field_mode(&api(401)));
assert!(!is_recoverable_in_field_mode(&api(403)));
assert!(!is_recoverable_in_field_mode(&api(422)));
assert!(!is_recoverable_in_field_mode(&Error::Validation(
"x".into()
)));
}
#[test]
fn permanent_failure_classification() {
let api = |status: u16, code: Option<ErrorCode>| Error::Api {
status,
code,
message: "x".into(),
request_id: None,
retry_after: None,
details: None,
};
assert!(is_permanent_extend_failure(&api(
410,
Some(ErrorCode::ReservationExpired)
)));
assert!(is_permanent_extend_failure(&api(
409,
Some(ErrorCode::ReservationFinalized)
)));
assert!(is_permanent_extend_failure(&api(
409,
Some(ErrorCode::MaxExtensionsExceeded)
)));
assert!(is_permanent_extend_failure(&api(410, None)));
assert!(is_permanent_extend_failure(&api(
409,
Some(ErrorCode::TenantClosed)
)));
assert!(is_permanent_extend_failure(&api(
404,
Some(ErrorCode::NotFound)
)));
assert!(is_permanent_extend_failure(&api(404, None)));
assert!(!is_permanent_extend_failure(&api(
500,
Some(ErrorCode::InternalError)
)));
assert!(!is_permanent_extend_failure(&api(429, None)));
assert!(!is_permanent_extend_failure(&Error::Validation("x".into())));
}
}