use super::*;
use crate::testing::acquire_test_state_lock;
use std::sync::atomic::{AtomicU32, Ordering};
use std::sync::Arc;
#[derive(Debug, Clone, PartialEq)]
enum TestErr {
Timeout,
CircuitOpen,
Other,
}
#[test]
fn inner_invokes_timeout_err_when_wrapper_timeout_fires() {
let _guard = acquire_test_state_lock();
_reset_resilience_for_tests();
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("runtime");
let policy = ExporterPolicy {
retries: 0,
backoff_seconds: 0.0,
timeout_seconds: 0.05,
fail_open: false,
allow_blocking_in_event_loop: false,
};
let te_calls = Arc::new(AtomicU32::new(0));
let is_calls = Arc::new(AtomicU32::new(0));
let co_calls = Arc::new(AtomicU32::new(0));
let result = runtime.block_on(async {
run_with_resilience_inner::<_, _, (), TestErr>(
Signal::Logs,
&policy,
|| async {
tokio::time::sleep(Duration::from_millis(500)).await;
Ok(())
},
{
let c = Arc::clone(&te_calls);
move |_dur| {
c.fetch_add(1, Ordering::SeqCst);
TestErr::Timeout
}
},
{
let c = Arc::clone(&is_calls);
move |_e| {
c.fetch_add(1, Ordering::SeqCst);
false
}
},
{
let c = Arc::clone(&co_calls);
move || {
c.fetch_add(1, Ordering::SeqCst);
TestErr::CircuitOpen
}
},
)
.await
});
assert_eq!(result, Err(TestErr::Timeout));
assert_eq!(
te_calls.load(Ordering::SeqCst),
1,
"timeout_err called exactly once on wrapper-timeout fire"
);
assert_eq!(
is_calls.load(Ordering::SeqCst),
0,
"is_sdk_timeout must NOT be called when wrapper_timeout is already true (short-circuit)"
);
assert_eq!(
co_calls.load(Ordering::SeqCst),
0,
"circuit_open_err must not fire when the circuit is closed"
);
}
#[test]
fn inner_invokes_is_sdk_timeout_for_each_operation_error() {
let _guard = acquire_test_state_lock();
_reset_resilience_for_tests();
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("runtime");
let policy = ExporterPolicy {
retries: 2,
backoff_seconds: 0.0,
timeout_seconds: 0.0, fail_open: false,
allow_blocking_in_event_loop: false,
};
let is_calls = Arc::new(AtomicU32::new(0));
let result = runtime.block_on(async {
run_with_resilience_inner::<_, _, (), TestErr>(
Signal::Logs,
&policy,
|| async { Err::<(), TestErr>(TestErr::Other) },
|_dur| TestErr::Timeout,
{
let c = Arc::clone(&is_calls);
move |_e| {
c.fetch_add(1, Ordering::SeqCst);
false
}
},
|| TestErr::CircuitOpen,
)
.await
});
assert_eq!(result, Err(TestErr::Other));
assert_eq!(
is_calls.load(Ordering::SeqCst),
3,
"is_sdk_timeout invoked once per failed attempt (retries+1 total)"
);
}
#[test]
fn inner_invokes_circuit_open_err_when_breaker_open_and_fail_closed() {
let _guard = acquire_test_state_lock();
_reset_resilience_for_tests();
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("runtime");
let policy = ExporterPolicy {
retries: 0,
backoff_seconds: 0.0,
timeout_seconds: 1.0, fail_open: false,
allow_blocking_in_event_loop: false,
};
{
let mut lock = crate::_lock::lock(circuits());
let state = lock
.get_mut(&Signal::Logs)
.expect("logs state should exist");
state.consecutive_timeouts = CIRCUIT_BREAKER_THRESHOLD;
state.open_count = 1;
state.tripped_at = Some(Instant::now());
}
let co_calls = Arc::new(AtomicU32::new(0));
let op_calls = Arc::new(AtomicU32::new(0));
let op_clone = Arc::clone(&op_calls);
let result = runtime.block_on(async {
run_with_resilience_inner::<_, _, (), TestErr>(
Signal::Logs,
&policy,
move || {
let c = Arc::clone(&op_clone);
async move {
c.fetch_add(1, Ordering::SeqCst);
Ok(())
}
},
|_dur| TestErr::Timeout,
|_e| false,
{
let c = Arc::clone(&co_calls);
move || {
c.fetch_add(1, Ordering::SeqCst);
TestErr::CircuitOpen
}
},
)
.await
});
assert_eq!(result, Err(TestErr::CircuitOpen));
assert_eq!(
co_calls.load(Ordering::SeqCst),
1,
"circuit_open_err invoked exactly once when breaker rejects the call"
);
assert_eq!(
op_calls.load(Ordering::SeqCst),
0,
"operation must not run when the breaker is open"
);
}
#[test]
fn inner_returns_none_when_breaker_open_and_fail_open() {
let _guard = acquire_test_state_lock();
_reset_resilience_for_tests();
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("runtime");
{
let mut lock = crate::_lock::lock(circuits());
let state = lock
.get_mut(&Signal::Logs)
.expect("logs state should exist");
state.consecutive_timeouts = CIRCUIT_BREAKER_THRESHOLD;
state.tripped_at = Some(Instant::now());
}
let result = runtime.block_on(async {
run_with_resilience_inner::<_, _, (), TestErr>(
Signal::Logs,
&ExporterPolicy {
retries: 0,
backoff_seconds: 0.0,
timeout_seconds: 1.0,
fail_open: true,
allow_blocking_in_event_loop: false,
},
|| async { Ok(()) },
|_dur| TestErr::Timeout,
|_err| false,
|| TestErr::CircuitOpen,
)
.await
});
assert_eq!(result, Ok(None));
}
#[test]
fn inner_success_path_covers_timeout_bypass_and_success_callbacks() {
let _guard = acquire_test_state_lock();
_reset_resilience_for_tests();
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("runtime");
let result = runtime.block_on(async {
run_with_resilience_inner::<_, _, u32, TestErr>(
Signal::Logs,
&ExporterPolicy {
retries: 0,
backoff_seconds: 0.0,
timeout_seconds: 0.0,
fail_open: false,
allow_blocking_in_event_loop: false,
},
|| async { Ok(7_u32) },
|_dur| TestErr::Timeout,
|_err| false,
|| TestErr::CircuitOpen,
)
.await
});
assert_eq!(result, Ok(Some(7_u32)));
let state = crate::_lock::lock(circuits());
assert_eq!(
state
.get(&Signal::Logs)
.expect("logs state should exist")
.consecutive_timeouts,
0
);
}
#[test]
fn inner_retries_then_succeeds_after_non_timeout_error() {
let _guard = acquire_test_state_lock();
_reset_resilience_for_tests();
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("runtime");
let attempts = Arc::new(AtomicU32::new(0));
let timeout_calls = Arc::new(AtomicU32::new(0));
let is_timeout_calls = Arc::new(AtomicU32::new(0));
let result = runtime.block_on(async {
run_with_resilience_inner::<_, _, &'static str, TestErr>(
Signal::Logs,
&ExporterPolicy {
retries: 1,
backoff_seconds: 0.001,
timeout_seconds: 1.0,
fail_open: false,
allow_blocking_in_event_loop: false,
},
{
let attempts = Arc::clone(&attempts);
move || {
let attempts = Arc::clone(&attempts);
async move {
if attempts.fetch_add(1, Ordering::SeqCst) == 0 {
Err(TestErr::Other)
} else {
Ok("ok")
}
}
}
},
{
let timeout_calls = Arc::clone(&timeout_calls);
move |_dur| {
timeout_calls.fetch_add(1, Ordering::SeqCst);
TestErr::Timeout
}
},
{
let is_timeout_calls = Arc::clone(&is_timeout_calls);
move |_err| {
is_timeout_calls.fetch_add(1, Ordering::SeqCst);
false
}
},
|| TestErr::CircuitOpen,
)
.await
});
assert_eq!(result, Ok(Some("ok")));
assert_eq!(attempts.load(Ordering::SeqCst), 2);
assert_eq!(timeout_calls.load(Ordering::SeqCst), 0);
assert_eq!(is_timeout_calls.load(Ordering::SeqCst), 1);
}
#[test]
fn inner_half_open_probe_failure_reopens_circuit() {
let _guard = acquire_test_state_lock();
_reset_resilience_for_tests();
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("runtime");
{
let mut lock = crate::_lock::lock(circuits());
let state = lock
.get_mut(&Signal::Logs)
.expect("logs state should exist");
state.consecutive_timeouts = CIRCUIT_BREAKER_THRESHOLD;
state.open_count = 1;
state.tripped_at = Some(Instant::now() - CIRCUIT_COOLDOWN - Duration::from_secs(1));
}
let result = runtime.block_on(async {
run_with_resilience_inner::<_, _, (), TestErr>(
Signal::Logs,
&ExporterPolicy {
retries: 0,
backoff_seconds: 0.0,
timeout_seconds: 1.0,
fail_open: true,
allow_blocking_in_event_loop: false,
},
|| async { Err(TestErr::Other) },
|_dur| TestErr::Timeout,
|_err| true,
|| TestErr::CircuitOpen,
)
.await
});
assert_eq!(result, Ok(None));
let state = crate::_lock::lock(circuits());
let logs = state.get(&Signal::Logs).expect("logs state should exist");
assert!(!logs.half_open_probing);
assert!(logs.open_count >= 2);
assert!(logs.tripped_at.is_some());
}