use std::sync::atomic::{AtomicU64, Ordering};
use std::time::Duration;
use runcycles::models::*;
use runcycles::CyclesClient;
use serde_json::json;
use wiremock::matchers::{method, path};
use wiremock::{Mock, MockServer, Request, Respond, ResponseTemplate};
const BASE_MS: u64 = 1_700_000_000_000;
const MARGIN: Duration = Duration::from_millis(450);
fn extend_path(rsv_id: &str) -> String {
format!("/v1/reservations/{rsv_id}/extend")
}
fn initial_expiry(granted_ttl_ms: u64) -> u64 {
BASE_MS + granted_ttl_ms
}
struct ExpirySequence {
base: u64,
step: u64,
first_status: &'static str,
rest_status: &'static str,
remaining_ttl_ms: Option<u64>,
remaining_on_first_n: Option<u64>,
remaining_first: Option<u64>,
calls: AtomicU64,
}
impl ExpirySequence {
fn granting(base: u64, step: u64) -> Self {
Self {
base,
step,
first_status: "ACTIVE",
rest_status: "ACTIVE",
remaining_ttl_ms: None,
remaining_on_first_n: None,
remaining_first: None,
calls: AtomicU64::new(0),
}
}
fn with_remaining(mut self, remaining_ttl_ms: u64) -> Self {
self.remaining_ttl_ms = Some(remaining_ttl_ms);
self
}
fn remaining_only_on_first(mut self, n: u64) -> Self {
self.remaining_on_first_n = Some(n);
self
}
fn with_first_remaining(mut self, remaining_ttl_ms: u64) -> Self {
self.remaining_first = Some(remaining_ttl_ms);
self
}
}
impl Respond for ExpirySequence {
fn respond(&self, _request: &Request) -> ResponseTemplate {
let n = self.calls.fetch_add(1, Ordering::SeqCst) + 1;
let status = if n == 1 {
self.first_status
} else {
self.rest_status
};
let mut body = json!({
"status": status,
"expires_at_ms": self.base + n * self.step
});
let remaining = if n == 1 {
self.remaining_first.or(self.remaining_ttl_ms)
} else {
self.remaining_ttl_ms
};
if let Some(remaining) = remaining {
if self.remaining_on_first_n.is_none_or(|first_n| n <= first_n) {
body["remaining_ttl_ms"] = json!(remaining);
}
}
ResponseTemplate::new(200).set_body_json(body)
}
}
struct LeadClampEcho {
base: u64,
created: std::time::Instant,
remaining_ttl_ms: Option<u64>,
}
impl LeadClampEcho {
fn holding(base: u64) -> Self {
Self {
base,
created: std::time::Instant::now(),
remaining_ttl_ms: None,
}
}
fn with_remaining(mut self, remaining_ttl_ms: u64) -> Self {
self.remaining_ttl_ms = Some(remaining_ttl_ms);
self
}
}
impl Respond for LeadClampEcho {
fn respond(&self, _request: &Request) -> ResponseTemplate {
let elapsed = u64::try_from(self.created.elapsed().as_millis()).unwrap_or(u64::MAX);
let mut body = json!({
"status": "ACTIVE",
"expires_at_ms": self.base + elapsed
});
if let Some(remaining) = self.remaining_ttl_ms {
body["remaining_ttl_ms"] = json!(remaining);
}
ResponseTemplate::new(200).set_body_json(body)
}
}
async fn mount_reserve_full(
server: &MockServer,
rsv_id: &str,
expires_at_ms: u64,
remaining_ttl_ms: Option<u64>,
) {
let mut body = json!({
"decision": "ALLOW",
"reservation_id": rsv_id,
"affected_scopes": ["tenant:acme"],
"expires_at_ms": expires_at_ms
});
if let Some(remaining) = remaining_ttl_ms {
body["remaining_ttl_ms"] = json!(remaining);
}
Mock::given(method("POST"))
.and(path("/v1/reservations"))
.respond_with(ResponseTemplate::new(200).set_body_json(body))
.mount(server)
.await;
Mock::given(method("POST"))
.and(path(format!("/v1/reservations/{rsv_id}/release")))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({
"status": "RELEASED",
"released": {"unit": "USD_MICROCENTS", "amount": 5000}
})))
.mount(server)
.await;
}
async fn mount_reserve(server: &MockServer, rsv_id: &str, expires_at_ms: u64) {
mount_reserve_full(server, rsv_id, expires_at_ms, None).await;
}
async fn do_reserve(server: &MockServer, requested_ttl_ms: u64) -> runcycles::ReservationGuard {
let client = CyclesClient::builder("key", server.uri()).build();
let req = ReservationCreateRequest::builder()
.subject(Subject {
tenant: Some("acme".into()),
..Default::default()
})
.action(Action::new("llm.completion", "gpt-4o"))
.estimate(Amount::usd_microcents(5000))
.ttl_ms(requested_ttl_ms)
.build();
client.reserve(req).await.unwrap()
}
async fn reserve_with_ttl(
server: &MockServer,
rsv_id: &str,
ttl_ms: u64,
) -> runcycles::ReservationGuard {
mount_reserve(server, rsv_id, initial_expiry(ttl_ms)).await;
do_reserve(server, ttl_ms).await
}
async fn extend_bodies(
server: &MockServer,
rsv_id: &str,
requested_ttl_ms: u64,
) -> Vec<serde_json::Value> {
let requests = server.received_requests().await.unwrap();
requests
.iter()
.filter(|r| r.url.path() == extend_path(rsv_id))
.map(|r| {
let body: serde_json::Value = r.body_json().unwrap();
assert_eq!(
body["extend_by_ms"], requested_ttl_ms,
"extend amount must be the requested ttl"
);
body
})
.collect()
}
async fn extend_calls(server: &MockServer, rsv_id: &str, requested_ttl_ms: u64) -> usize {
extend_bodies(server, rsv_id, requested_ttl_ms).await.len()
}
#[tokio::test]
async fn heartbeat_first_beat_immediate_then_lead_bound_cadence() {
let server = MockServer::start().await;
let rsv_id = "rsv_hb_lead";
const TTL: u64 = 2_000;
const BEAT: Duration = Duration::from_millis(1_000);
Mock::given(method("POST"))
.and(path(extend_path(rsv_id)))
.respond_with(ExpirySequence::granting(initial_expiry(TTL), TTL))
.mount(&server)
.await;
let guard = reserve_with_ttl(&server, rsv_id, TTL).await;
tokio::time::sleep(MARGIN).await;
assert_eq!(
extend_calls(&server, rsv_id, TTL).await,
1,
"the first extend must fire immediately — well before ttl/2 = 1000 ms \
(any bounded delay can outlive a small tenant-capped lease)"
);
tokio::time::sleep(BEAT).await;
assert_eq!(
extend_calls(&server, rsv_id, TTL).await,
2,
"beat 2 at 1000 ms must extend (lead_min 1000 < 3000)"
);
tokio::time::sleep(BEAT).await;
assert_eq!(
extend_calls(&server, rsv_id, TTL).await,
3,
"beat 3 at 2000 ms must extend (lead_min 2000 < 3000)"
);
tokio::time::sleep(BEAT).await;
assert_eq!(
extend_calls(&server, rsv_id, TTL).await,
3,
"beat 4 at 3000 ms must skip (lead_min 3000 = 1.5·grant, inclusive)"
);
tokio::time::sleep(BEAT).await;
assert_eq!(
extend_calls(&server, rsv_id, TTL).await,
4,
"beat 5 at 4000 ms must extend again (lead_min back to 2000)"
);
guard.release("test done").await.unwrap();
}
#[tokio::test]
async fn heartbeat_capped_grant_derives_cadence_from_observed_grants() {
let server = MockServer::start().await;
let rsv_id = "rsv_hb_capped";
const REQUESTED: u64 = 8_000;
const GRANTED: u64 = 2_000;
const BEAT: Duration = Duration::from_millis(1_000);
mount_reserve(&server, rsv_id, initial_expiry(GRANTED)).await;
Mock::given(method("POST"))
.and(path(extend_path(rsv_id)))
.respond_with(ExpirySequence::granting(initial_expiry(GRANTED), GRANTED))
.mount(&server)
.await;
let guard = do_reserve(&server, REQUESTED).await;
tokio::time::sleep(MARGIN).await;
assert_eq!(
extend_calls(&server, rsv_id, REQUESTED).await,
1,
"the first beat must be immediate — the only schedule that cannot \
outlive the silently-capped 2000 ms lease"
);
tokio::time::sleep(BEAT).await;
assert_eq!(
extend_calls(&server, rsv_id, REQUESTED).await,
2,
"beat 2 at 1000 ms must extend (cadence follows the observed grant)"
);
tokio::time::sleep(BEAT).await;
assert_eq!(
extend_calls(&server, rsv_id, REQUESTED).await,
3,
"beat 3 at 2000 ms must extend (lead_min 2000 < 1.5·grant)"
);
guard.release("test done").await.unwrap();
}
#[tokio::test]
async fn heartbeat_first_beat_failure_retries_at_held_cadence_with_same_key() {
let server = MockServer::start().await;
let rsv_id = "rsv_hb_key";
const TTL: u64 = 2_000;
const BEAT: Duration = Duration::from_millis(1_000);
Mock::given(method("POST"))
.and(path(extend_path(rsv_id)))
.respond_with(ResponseTemplate::new(503).set_body_json(json!({
"error": "INTERNAL_ERROR",
"message": "boom"
})))
.up_to_n_times(1)
.mount(&server)
.await;
Mock::given(method("POST"))
.and(path(extend_path(rsv_id)))
.respond_with(ExpirySequence::granting(initial_expiry(TTL), TTL))
.mount(&server)
.await;
let guard = reserve_with_ttl(&server, rsv_id, TTL).await;
tokio::time::sleep(MARGIN).await;
assert_eq!(
extend_calls(&server, rsv_id, TTL).await,
1,
"the failed immediate beat must wait the held cadence before \
retrying — never retry at zero delay"
);
tokio::time::sleep(BEAT).await;
let bodies = extend_bodies(&server, rsv_id, TTL).await;
assert_eq!(
bodies.len(),
2,
"beat 2 at 1000 ms must retry the failed extend"
);
assert_eq!(
bodies[0]["idempotency_key"], bodies[1]["idempotency_key"],
"the retry must reuse the failed attempt's idempotency key \
(dedupe against an applied-but-lost extension)"
);
tokio::time::sleep(BEAT).await;
let bodies = extend_bodies(&server, rsv_id, TTL).await;
assert_eq!(bodies.len(), 3, "beat 3 must extend (lead_min 0 < 3000)");
assert_ne!(
bodies[2]["idempotency_key"], bodies[0]["idempotency_key"],
"after a success the next extend must use a fresh idempotency key"
);
guard.release("test done").await.unwrap();
}
async fn assert_permanent_stop(rsv_id: &str, status: u16, error_code: &str) {
let server = MockServer::start().await;
const TTL: u64 = 2_000;
Mock::given(method("POST"))
.and(path(extend_path(rsv_id)))
.respond_with(ResponseTemplate::new(status).set_body_json(json!({
"error": error_code,
"message": "permanent"
})))
.mount(&server)
.await;
let guard = reserve_with_ttl(&server, rsv_id, TTL).await;
tokio::time::sleep(MARGIN).await;
assert_eq!(
extend_calls(&server, rsv_id, TTL).await,
1,
"the immediate beat attempts and hits the permanent failure ({error_code})"
);
tokio::time::sleep(Duration::from_millis(2_000)).await;
assert_eq!(
extend_calls(&server, rsv_id, TTL).await,
1,
"the heartbeat must stop after {error_code} — no retries"
);
guard.release("test done").await.unwrap();
}
#[tokio::test]
async fn heartbeat_stops_on_max_extensions_exceeded() {
assert_permanent_stop("rsv_hb_perm", 409, "MAX_EXTENSIONS_EXCEEDED").await;
}
#[tokio::test]
async fn heartbeat_stops_on_tenant_closed() {
assert_permanent_stop("rsv_hb_closed", 409, "TENANT_CLOSED").await;
}
#[tokio::test]
async fn heartbeat_stops_on_not_found() {
assert_permanent_stop("rsv_hb_404", 404, "NOT_FOUND").await;
}
#[tokio::test]
async fn heartbeat_small_ttl_stays_alive() {
let server = MockServer::start().await;
let rsv_id = "rsv_hb_small";
const TTL: u64 = 1_200;
const SMALL_MARGIN: Duration = Duration::from_millis(250);
const SMALL_BEAT: Duration = Duration::from_millis(600);
Mock::given(method("POST"))
.and(path(extend_path(rsv_id)))
.respond_with(ExpirySequence::granting(initial_expiry(TTL), TTL))
.mount(&server)
.await;
let guard = reserve_with_ttl(&server, rsv_id, TTL).await;
tokio::time::sleep(SMALL_MARGIN).await;
assert_eq!(
extend_calls(&server, rsv_id, TTL).await,
1,
"the first beat must be immediate"
);
tokio::time::sleep(SMALL_BEAT).await;
assert_eq!(
extend_calls(&server, rsv_id, TTL).await,
2,
"beat 2 at 600 ms must extend (no 1-second floor; lead_min 600 < 1800)"
);
tokio::time::sleep(SMALL_BEAT).await;
assert_eq!(
extend_calls(&server, rsv_id, TTL).await,
3,
"beat 3 at 1200 ms must extend (lead_min 1200 < 1800)"
);
tokio::time::sleep(SMALL_BEAT).await;
assert_eq!(
extend_calls(&server, rsv_id, TTL).await,
3,
"beat 4 at 1800 ms skips once the lead lower bound reaches 1.5·grant"
);
tokio::time::sleep(SMALL_BEAT).await;
assert_eq!(
extend_calls(&server, rsv_id, TTL).await,
4,
"beat 5 at 2400 ms must extend again"
);
guard.release("test done").await.unwrap();
}
#[tokio::test]
async fn heartbeat_clamped_grant_still_tightens_cadence() {
let server = MockServer::start().await;
let rsv_id = "rsv_hb_clamp";
const TTL: u64 = 4_000;
Mock::given(method("POST"))
.and(path(extend_path(rsv_id)))
.respond_with(ExpirySequence::granting(initial_expiry(TTL), TTL / 4))
.mount(&server)
.await;
let guard = reserve_with_ttl(&server, rsv_id, TTL).await;
tokio::time::sleep(Duration::from_millis(250)).await;
assert_eq!(
extend_calls(&server, rsv_id, TTL).await,
1,
"the first beat must be immediate"
);
tokio::time::sleep(Duration::from_millis(1_000)).await;
assert_eq!(
extend_calls(&server, rsv_id, TTL).await,
3,
"clamped per-extend grants must tighten the cadence to grant/2 = \
500 ms (beats at 0/500/1000 ms) — not hold it at requested/2"
);
guard.release("test done").await.unwrap();
}
#[tokio::test]
async fn heartbeat_lead_clamp_holds_cadence_instead_of_collapsing() {
let server = MockServer::start().await;
let rsv_id = "rsv_hb_leadclamp";
const TTL: u64 = 3_000;
mount_reserve(&server, rsv_id, initial_expiry(TTL)).await;
Mock::given(method("POST"))
.and(path(extend_path(rsv_id)))
.respond_with(LeadClampEcho::holding(initial_expiry(TTL)))
.mount(&server)
.await;
let guard = do_reserve(&server, TTL).await;
tokio::time::sleep(Duration::from_millis(3_250)).await;
let after_3s = extend_calls(&server, rsv_id, TTL).await;
assert!(
(3..=4).contains(&after_3s),
"lead-clamp cadence must hold at 1500 ms (3-4 extends by 3.25 s); \
a floor-collapsed cadence would have sent ~7, got {after_3s}"
);
tokio::time::sleep(Duration::from_millis(1_400)).await;
let after_4s = extend_calls(&server, rsv_id, TTL).await;
assert!(
(4..=5).contains(&after_4s),
"lead-clamp cadence must stay held (4-5 extends by 4.65 s); a \
floor-collapsed cadence would have sent ~9, got {after_4s}"
);
assert!(
after_4s > after_3s,
"the heartbeat must keep extending at the held cadence (liveness)"
);
guard.release("test done").await.unwrap();
}
#[tokio::test]
async fn heartbeat_zero_grant_immediate_prime_holds_cadence() {
let server = MockServer::start().await;
let rsv_id = "rsv_hb_zerogrant";
const TTL: u64 = 2_000;
mount_reserve(&server, rsv_id, initial_expiry(TTL)).await;
Mock::given(method("POST"))
.and(path(extend_path(rsv_id)))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({
"status": "ACTIVE",
"expires_at_ms": initial_expiry(TTL)
})))
.mount(&server)
.await;
let guard = do_reserve(&server, TTL).await;
tokio::time::sleep(MARGIN).await;
assert_eq!(
extend_calls(&server, rsv_id, TTL).await,
1,
"the first beat is immediate and observes a zero grant"
);
tokio::time::sleep(Duration::from_millis(2_000)).await;
assert_eq!(
extend_calls(&server, rsv_id, TTL).await,
3,
"zero grants must hold the cadence at min(requested/2, 30 s), \
never collapse to the 500 ms floor"
);
guard.release("test done").await.unwrap();
}
#[tokio::test]
async fn heartbeat_fallback_unknown_status_retries_with_same_key() {
let server = MockServer::start().await;
let rsv_id = "rsv_hb_unknown";
const TTL: u64 = 2_000;
const BEAT: Duration = Duration::from_millis(1_000);
Mock::given(method("POST"))
.and(path(extend_path(rsv_id)))
.respond_with(ExpirySequence {
first_status: "SUSPENDED",
..ExpirySequence::granting(initial_expiry(TTL), TTL)
})
.mount(&server)
.await;
let guard = reserve_with_ttl(&server, rsv_id, TTL).await;
tokio::time::sleep(MARGIN).await;
assert_eq!(
extend_calls(&server, rsv_id, TTL).await,
1,
"the immediate beat extends and receives the unknown-status 200"
);
tokio::time::sleep(BEAT + MARGIN).await;
let bodies = extend_bodies(&server, rsv_id, TTL).await;
assert_eq!(bodies.len(), 2, "beat 2 extends (lead_min 1000 < 3000)");
assert_eq!(
bodies[0]["idempotency_key"], bodies[1]["idempotency_key"],
"the unknown-status 200 is ambiguous: beat 2 must retry the same key"
);
guard.release("test done").await.unwrap();
}
#[tokio::test]
async fn heartbeat_transport_timeout_is_observable_nonfatal_and_same_key() {
let server = MockServer::start().await;
let rsv_id = "rsv_hb_transport";
const TTL: u64 = 2_000;
mount_reserve(&server, rsv_id, initial_expiry(TTL)).await;
Mock::given(method("POST"))
.and(path(extend_path(rsv_id)))
.respond_with(
ResponseTemplate::new(200)
.set_delay(Duration::from_millis(250))
.set_body_json(json!({
"status": "ACTIVE",
"expires_at_ms": initial_expiry(TTL) + TTL
})),
)
.mount(&server)
.await;
Mock::given(method("POST"))
.and(path(format!("/v1/reservations/{rsv_id}/commit")))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({
"status": "COMMITTED",
"charged": {"unit": "USD_MICROCENTS", "amount": 5000}
})))
.mount(&server)
.await;
let client = CyclesClient::builder("key", server.uri())
.connect_timeout(Duration::from_millis(50))
.read_timeout(Duration::from_millis(50))
.build();
let request = ReservationCreateRequest::builder()
.subject(Subject {
tenant: Some("acme".into()),
..Default::default()
})
.action(Action::new("llm.completion", "gpt-4o"))
.estimate(Amount::usd_microcents(5000))
.ttl_ms(TTL)
.build();
let guard = client.reserve(request).await.unwrap();
tokio::time::sleep(Duration::from_millis(1_450)).await;
let bodies = extend_bodies(&server, rsv_id, TTL).await;
assert!(bodies.len() >= 2);
assert_eq!(bodies[0]["idempotency_key"], bodies[1]["idempotency_key"]);
guard
.commit(
CommitRequest::builder()
.idempotency_key(IdempotencyKey::new("commit-after-heartbeat"))
.actual(Amount::usd_microcents(5000))
.build(),
)
.await
.unwrap();
let commit_count = server
.received_requests()
.await
.unwrap()
.iter()
.filter(|request| request.url.path().ends_with("/commit"))
.count();
assert_eq!(commit_count, 1);
}
async fn do_reserve_fast(
server: &MockServer,
requested_ttl_ms: u64,
) -> runcycles::ReservationGuard {
let client = CyclesClient::builder("key", server.uri())
.connect_timeout(Duration::from_millis(500))
.read_timeout(Duration::from_millis(500))
.build();
let req = ReservationCreateRequest::builder()
.subject(Subject {
tenant: Some("acme".into()),
..Default::default()
})
.action(Action::new("llm.completion", "gpt-4o"))
.estimate(Amount::usd_microcents(5000))
.ttl_ms(requested_ttl_ms)
.build();
client.reserve(req).await.unwrap()
}
#[tokio::test]
async fn heartbeat_remaining_ttl_schedules_normatively_and_bypasses_skip() {
let server = MockServer::start().await;
let rsv_id = "rsv_hb_norm";
const REQUESTED: u64 = 8_000;
const REMAINING: u64 = 4_000;
const BEAT: Duration = Duration::from_millis(1_000);
mount_reserve_full(&server, rsv_id, initial_expiry(REMAINING), Some(REMAINING)).await;
Mock::given(method("POST"))
.and(path(extend_path(rsv_id)))
.respond_with(
ExpirySequence::granting(initial_expiry(REMAINING), REMAINING)
.with_remaining(REMAINING),
)
.mount(&server)
.await;
let guard = do_reserve_fast(&server, REQUESTED).await;
tokio::time::sleep(MARGIN).await;
assert_eq!(
extend_calls(&server, rsv_id, REQUESTED).await,
0,
"with remaining_ttl_ms on the create response there is no immediate \
primed extension — the first beat is scheduled normatively at \
lead_floor − retry_reserve ≈ 1000 ms"
);
tokio::time::sleep(BEAT).await;
assert_eq!(
extend_calls(&server, rsv_id, REQUESTED).await,
1,
"beat 1 lands ≈ 1000 ms — inside the real 4 s lease (a requested/2 \
= 4 s schedule would have been at its very end)"
);
tokio::time::sleep(BEAT + BEAT).await;
assert_eq!(
extend_calls(&server, rsv_id, REQUESTED).await,
3,
"normative beats every ≈ 1000 ms"
);
tokio::time::sleep(BEAT).await;
assert_eq!(
extend_calls(&server, rsv_id, REQUESTED).await,
4,
"beat 4 must extend — the lead_min heuristic (which would skip here \
on the accumulated +4000 grants) must be bypassed in field mode"
);
guard.release("test done").await.unwrap();
}
#[tokio::test]
async fn heartbeat_remaining_ttl_zero_delay_guard_stops_after_one_fresh_attempt() {
let server = MockServer::start().await;
let rsv_id = "rsv_hb_norm_zero";
const REQUESTED: u64 = 8_000;
const REMAINING: u64 = 2_500;
mount_reserve_full(&server, rsv_id, initial_expiry(REMAINING), Some(REMAINING)).await;
Mock::given(method("POST"))
.and(path(extend_path(rsv_id)))
.respond_with(
ExpirySequence::granting(initial_expiry(REMAINING), REMAINING)
.with_remaining(REMAINING),
)
.mount(&server)
.await;
let guard = do_reserve_fast(&server, REQUESTED).await;
tokio::time::sleep(MARGIN).await;
assert_eq!(
extend_calls(&server, rsv_id, REQUESTED).await,
1,
"next_delay = 0 from the create permits exactly one immediate \
fresh-key extension"
);
tokio::time::sleep(Duration::from_millis(2_000)).await;
assert_eq!(
extend_calls(&server, rsv_id, REQUESTED).await,
1,
"a second zero-delay success must stop the heartbeat (lease shorter \
than the retry-safety budget) — no tight extension loop"
);
guard.release("test done").await.unwrap();
}
#[tokio::test]
async fn heartbeat_remaining_ttl_under_max_lead_clamp() {
let server = MockServer::start().await;
let rsv_id = "rsv_hb_norm_clamp";
const TTL: u64 = 5_000;
mount_reserve_full(&server, rsv_id, initial_expiry(TTL), Some(TTL)).await;
Mock::given(method("POST"))
.and(path(extend_path(rsv_id)))
.respond_with(LeadClampEcho::holding(initial_expiry(TTL)).with_remaining(TTL))
.mount(&server)
.await;
let guard = do_reserve_fast(&server, TTL).await;
tokio::time::sleep(MARGIN).await;
assert_eq!(
extend_calls(&server, rsv_id, TTL).await,
0,
"no primed extension under a max-lead clamp when the create reports \
remaining_ttl_ms — the first beat waits ≈ 2000 ms"
);
tokio::time::sleep(Duration::from_millis(2_000)).await;
assert_eq!(
extend_calls(&server, rsv_id, TTL).await,
1,
"beat 1 at remaining − retry_reserve ≈ 2000 ms"
);
tokio::time::sleep(Duration::from_millis(2_000)).await;
assert_eq!(
extend_calls(&server, rsv_id, TTL).await,
2,
"steady ≈ 2000 ms cadence from the reported remaining lease — never \
a collapse to the 500 ms floor"
);
guard.release("test done").await.unwrap();
}
#[tokio::test]
async fn heartbeat_remaining_ttl_disappearing_falls_back_to_heuristic() {
let server = MockServer::start().await;
let rsv_id = "rsv_hb_norm_gone";
const REQUESTED: u64 = 1_200;
const REMAINING: u64 = 4_000;
mount_reserve_full(&server, rsv_id, initial_expiry(REMAINING), Some(REMAINING)).await;
Mock::given(method("POST"))
.and(path(extend_path(rsv_id)))
.respond_with(
ExpirySequence::granting(initial_expiry(REMAINING), REQUESTED)
.with_remaining(REMAINING)
.remaining_only_on_first(1),
)
.mount(&server)
.await;
let guard = do_reserve_fast(&server, REQUESTED).await;
tokio::time::sleep(Duration::from_millis(250)).await;
assert_eq!(
extend_calls(&server, rsv_id, REQUESTED).await,
0,
"normative first beat (≈ 1000 ms): no immediate prime"
);
tokio::time::sleep(Duration::from_millis(1_200)).await;
assert_eq!(
extend_calls(&server, rsv_id, REQUESTED).await,
1,
"beat 1 lands normatively at ≈ 1000 ms"
);
tokio::time::sleep(Duration::from_millis(2_000)).await;
assert_eq!(
extend_calls(&server, rsv_id, REQUESTED).await,
4,
"after the field disappears the heuristic resumes at the observed \
grant/2 = 600 ms cadence (beats ≈ 2000/2600/3200 ms)"
);
guard.release("test done").await.unwrap();
}
#[tokio::test]
async fn heartbeat_remaining_ttl_transient_failure_retries_same_key() {
let server = MockServer::start().await;
let rsv_id = "rsv_hb_norm_retry";
const REQUESTED: u64 = 8_000;
const REMAINING: u64 = 5_000;
mount_reserve_full(&server, rsv_id, initial_expiry(REMAINING), Some(REMAINING)).await;
Mock::given(method("POST"))
.and(path(extend_path(rsv_id)))
.respond_with(ResponseTemplate::new(503).set_body_json(json!({
"error": "INTERNAL_ERROR",
"message": "boom"
})))
.up_to_n_times(1)
.mount(&server)
.await;
Mock::given(method("POST"))
.and(path(extend_path(rsv_id)))
.respond_with(
ExpirySequence::granting(initial_expiry(REMAINING), REMAINING)
.with_remaining(REMAINING),
)
.mount(&server)
.await;
let guard = do_reserve_fast(&server, REQUESTED).await;
tokio::time::sleep(Duration::from_millis(2_450)).await;
assert_eq!(
extend_calls(&server, rsv_id, REQUESTED).await,
1,
"the normative first beat at ≈ 2000 ms attempts and fails"
);
tokio::time::sleep(Duration::from_millis(1_000)).await;
let bodies = extend_bodies(&server, rsv_id, REQUESTED).await;
assert_eq!(
bodies.len(),
2,
"the retry lands min(30 s, lead/4, window) ≈ 750 ms after the failure"
);
assert_eq!(
bodies[0]["idempotency_key"], bodies[1]["idempotency_key"],
"the field-mode recovery retry must reuse the failed attempt's \
idempotency key"
);
guard.release("test done").await.unwrap();
}
#[tokio::test]
async fn heartbeat_remaining_ttl_ambiguous_2xx_recovers_with_same_key() {
let server = MockServer::start().await;
let rsv_id = "rsv_hb_norm_ambig";
const REQUESTED: u64 = 8_000;
const REMAINING: u64 = 5_000;
mount_reserve_full(&server, rsv_id, initial_expiry(REMAINING), Some(REMAINING)).await;
Mock::given(method("POST"))
.and(path(extend_path(rsv_id)))
.respond_with(ResponseTemplate::new(200).set_body_string("not json at all"))
.up_to_n_times(1)
.mount(&server)
.await;
Mock::given(method("POST"))
.and(path(extend_path(rsv_id)))
.respond_with(
ExpirySequence::granting(initial_expiry(REMAINING), REMAINING)
.with_remaining(REMAINING),
)
.mount(&server)
.await;
let guard = do_reserve_fast(&server, REQUESTED).await;
tokio::time::sleep(Duration::from_millis(2_450)).await;
assert_eq!(
extend_calls(&server, rsv_id, REQUESTED).await,
1,
"the ambiguous 200 must not count as success — one attempt so far, \
recovery pending"
);
tokio::time::sleep(Duration::from_millis(1_000)).await;
let bodies = extend_bodies(&server, rsv_id, REQUESTED).await;
assert_eq!(
bodies.len(),
2,
"the ambiguous 2xx is recovered like a transient failure ≈ 750 ms later"
);
assert_eq!(
bodies[0]["idempotency_key"], bodies[1]["idempotency_key"],
"ambiguous-2xx recovery must reuse the same idempotency key (the \
extension may have been applied)"
);
tokio::time::sleep(Duration::from_millis(1_800)).await;
let bodies = extend_bodies(&server, rsv_id, REQUESTED).await;
assert_eq!(bodies.len(), 3, "normative schedule resumes after recovery");
assert_ne!(
bodies[2]["idempotency_key"], bodies[0]["idempotency_key"],
"after a schema-valid success the next extend uses a fresh key"
);
guard.release("test done").await.unwrap();
}
#[tokio::test]
async fn heartbeat_remaining_ttl_recovery_stops_when_window_exhausted() {
let server = MockServer::start().await;
let rsv_id = "rsv_hb_norm_exhaust";
const REQUESTED: u64 = 8_000;
const REMAINING: u64 = 5_000;
mount_reserve_full(&server, rsv_id, initial_expiry(REMAINING), Some(REMAINING)).await;
Mock::given(method("POST"))
.and(path(extend_path(rsv_id)))
.respond_with(ResponseTemplate::new(500).set_body_json(json!({
"error": "INTERNAL_ERROR",
"message": "down"
})))
.mount(&server)
.await;
let guard = do_reserve_fast(&server, REQUESTED).await;
tokio::time::sleep(Duration::from_millis(4_200)).await;
let bodies = extend_bodies(&server, rsv_id, REQUESTED).await;
let n = bodies.len();
assert!(
(3..=5).contains(&n),
"recovery must retry while the window is non-negative and then stop \
(expected 3-5 attempts, got {n})"
);
for body in &bodies {
assert_eq!(
body["idempotency_key"], bodies[0]["idempotency_key"],
"every recovery retry must reuse the original idempotency key"
);
}
tokio::time::sleep(Duration::from_millis(1_000)).await;
assert_eq!(
extend_calls(&server, rsv_id, REQUESTED).await,
n,
"once retry_window went negative the heartbeat must stop for good"
);
guard.release("test done").await.unwrap();
}
#[tokio::test]
async fn heartbeat_remaining_ttl_429_retry_after_within_window() {
let server = MockServer::start().await;
let rsv_id = "rsv_hb_norm_429ok";
const REQUESTED: u64 = 8_000;
const REMAINING: u64 = 5_000;
mount_reserve_full(&server, rsv_id, initial_expiry(REMAINING), Some(REMAINING)).await;
Mock::given(method("POST"))
.and(path(extend_path(rsv_id)))
.respond_with(
ResponseTemplate::new(429)
.insert_header("retry-after", "0")
.set_body_json(json!({
"error": "LIMIT_EXCEEDED",
"message": "slow down"
})),
)
.up_to_n_times(1)
.mount(&server)
.await;
Mock::given(method("POST"))
.and(path(extend_path(rsv_id)))
.respond_with(
ExpirySequence::granting(initial_expiry(REMAINING), REMAINING)
.with_remaining(REMAINING),
)
.mount(&server)
.await;
let guard = do_reserve_fast(&server, REQUESTED).await;
tokio::time::sleep(Duration::from_millis(2_450)).await;
let bodies = extend_bodies(&server, rsv_id, REQUESTED).await;
assert_eq!(
bodies.len(),
2,
"the in-window 429 must be retried after exactly Retry-After (0 s → \
immediately)"
);
assert_eq!(
bodies[0]["idempotency_key"], bodies[1]["idempotency_key"],
"the 429 retry must reuse the same idempotency key"
);
guard.release("test done").await.unwrap();
}
#[tokio::test]
async fn heartbeat_remaining_ttl_429_retry_after_exceeding_window_stops() {
let server = MockServer::start().await;
let rsv_id = "rsv_hb_norm_429no";
const REQUESTED: u64 = 8_000;
const REMAINING: u64 = 5_000;
mount_reserve_full(&server, rsv_id, initial_expiry(REMAINING), Some(REMAINING)).await;
Mock::given(method("POST"))
.and(path(extend_path(rsv_id)))
.respond_with(
ResponseTemplate::new(429)
.insert_header("retry-after", "1")
.set_body_json(json!({
"error": "LIMIT_EXCEEDED",
"message": "slow down"
})),
)
.mount(&server)
.await;
let guard = do_reserve_fast(&server, REQUESTED).await;
tokio::time::sleep(Duration::from_millis(2_450)).await;
assert_eq!(
extend_calls(&server, rsv_id, REQUESTED).await,
1,
"the 429 at the first beat is received once"
);
tokio::time::sleep(Duration::from_millis(2_000)).await;
assert_eq!(
extend_calls(&server, rsv_id, REQUESTED).await,
1,
"Retry-After exceeding the retry window must stop the heartbeat — \
never an earlier retry that violates throttling"
);
guard.release("test done").await.unwrap();
}
#[tokio::test]
async fn heartbeat_remaining_ttl_other_4xx_stops_without_retry() {
let server = MockServer::start().await;
let rsv_id = "rsv_hb_norm_4xx";
const REQUESTED: u64 = 8_000;
const REMAINING: u64 = 5_000;
mount_reserve_full(&server, rsv_id, initial_expiry(REMAINING), Some(REMAINING)).await;
Mock::given(method("POST"))
.and(path(extend_path(rsv_id)))
.respond_with(ResponseTemplate::new(400).set_body_json(json!({
"error": "INVALID_REQUEST",
"message": "bad"
})))
.mount(&server)
.await;
let guard = do_reserve_fast(&server, REQUESTED).await;
tokio::time::sleep(Duration::from_millis(2_450)).await;
assert_eq!(
extend_calls(&server, rsv_id, REQUESTED).await,
1,
"the 400 at the first beat is received once"
);
tokio::time::sleep(Duration::from_millis(2_000)).await;
assert_eq!(
extend_calls(&server, rsv_id, REQUESTED).await,
1,
"another 4xx must stop the heartbeat without key rotation or retry"
);
guard.release("test done").await.unwrap();
}
#[tokio::test]
async fn heartbeat_remaining_ttl_non_200_2xx_is_ambiguous() {
let server = MockServer::start().await;
let rsv_id = "rsv_hb_norm_204";
const REQUESTED: u64 = 8_000;
const REMAINING: u64 = 5_000;
mount_reserve_full(&server, rsv_id, initial_expiry(REMAINING), Some(REMAINING)).await;
Mock::given(method("POST"))
.and(path(extend_path(rsv_id)))
.respond_with(ResponseTemplate::new(204))
.up_to_n_times(1)
.mount(&server)
.await;
Mock::given(method("POST"))
.and(path(extend_path(rsv_id)))
.respond_with(
ExpirySequence::granting(initial_expiry(REMAINING), REMAINING)
.with_remaining(REMAINING),
)
.mount(&server)
.await;
let guard = do_reserve_fast(&server, REQUESTED).await;
tokio::time::sleep(Duration::from_millis(2_450)).await;
assert_eq!(
extend_calls(&server, rsv_id, REQUESTED).await,
1,
"the 204 must not count as success — one attempt so far, recovery pending"
);
tokio::time::sleep(Duration::from_millis(1_000)).await;
let bodies = extend_bodies(&server, rsv_id, REQUESTED).await;
assert_eq!(
bodies.len(),
2,
"the ambiguous non-200 2xx is recovered like a transient failure"
);
assert_eq!(
bodies[0]["idempotency_key"], bodies[1]["idempotency_key"],
"ambiguous-2xx recovery must reuse the same idempotency key"
);
guard.release("test done").await.unwrap();
}
#[tokio::test]
async fn heartbeat_remaining_ttl_single_zero_delay_recovers() {
let server = MockServer::start().await;
let rsv_id = "rsv_hb_norm_zero1";
const REQUESTED: u64 = 8_000;
const REMAINING: u64 = 5_000;
mount_reserve_full(&server, rsv_id, initial_expiry(REMAINING), Some(REMAINING)).await;
Mock::given(method("POST"))
.and(path(extend_path(rsv_id)))
.respond_with(
ExpirySequence::granting(initial_expiry(REMAINING), REMAINING)
.with_remaining(REMAINING)
.with_first_remaining(2_500),
)
.mount(&server)
.await;
let guard = do_reserve_fast(&server, REQUESTED).await;
tokio::time::sleep(Duration::from_millis(2_450)).await;
let bodies = extend_bodies(&server, rsv_id, REQUESTED).await;
assert_eq!(
bodies.len(),
2,
"a single zero-delay success must be followed by exactly one \
immediate extension"
);
assert_ne!(
bodies[0]["idempotency_key"], bodies[1]["idempotency_key"],
"the immediate zero-delay extension is a FRESH attempt, not a \
same-key retry (the previous extend succeeded)"
);
tokio::time::sleep(Duration::from_millis(2_000)).await;
assert_eq!(
extend_calls(&server, rsv_id, REQUESTED).await,
3,
"a healthy lease after a single zero-delay success resumes the \
normative schedule (no stop, no tight loop)"
);
guard.release("test done").await.unwrap();
}
#[tokio::test]
async fn heartbeat_fallback_nonconformant_reserve_without_expiry() {
let server = MockServer::start().await;
let rsv_id = "rsv_hb_noexpiry";
const TTL: u64 = 2_000;
Mock::given(method("POST"))
.and(path("/v1/reservations"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({
"decision": "ALLOW",
"reservation_id": rsv_id,
"affected_scopes": ["tenant:acme"]
})))
.mount(&server)
.await;
Mock::given(method("POST"))
.and(path(format!("/v1/reservations/{rsv_id}/release")))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({
"status": "RELEASED",
"released": {"unit": "USD_MICROCENTS", "amount": 5000}
})))
.mount(&server)
.await;
Mock::given(method("POST"))
.and(path(extend_path(rsv_id)))
.respond_with(ExpirySequence::granting(initial_expiry(TTL), TTL))
.mount(&server)
.await;
let guard = do_reserve(&server, TTL).await;
tokio::time::sleep(MARGIN).await;
assert_eq!(
extend_calls(&server, rsv_id, TTL).await,
1,
"the fallback immediate prime fires even without an initial expiry \
sample (first grant falls back to the requested amount)"
);
tokio::time::sleep(Duration::from_millis(1_000)).await;
assert_eq!(
extend_calls(&server, rsv_id, TTL).await,
2,
"the requested-amount fallback grant paces the cadence at \
requested/2 = 1000 ms"
);
guard.release("test done").await.unwrap();
}