feat(nvim,telemetry): implement ADR-0102 circuit breaker and ADR-0103 knowledge graph projection
This commit is contained in:
1 parent
3d77e60a02
commit
ec977c24dd
7 files changed
+692
-11
No files matched your search
+145
-5
@@ -214,7 +214,98 @@ pub struct NvimRequest {
|
||||
pub reply: oneshot::Sender<Result<rmpv::Value, String>>,
|
||||
}
|
||||
|
||||
use std::sync::atomic::{AtomicU64, Ordering};
|
||||
use std::sync::atomic::{AtomicU32, AtomicU64, AtomicU8, Ordering};
|
||||
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub enum CircuitState {
|
||||
Closed = 0,
|
||||
Open = 1,
|
||||
HalfOpen = 2,
|
||||
}
|
||||
|
||||
pub struct NvimCircuitBreaker {
|
||||
state: AtomicU8,
|
||||
consecutive_failures: AtomicU32,
|
||||
last_state_change_millis: AtomicU64,
|
||||
failure_threshold: u32,
|
||||
cooldown_millis: u64,
|
||||
}
|
||||
|
||||
impl NvimCircuitBreaker {
|
||||
pub fn new(failure_threshold: u32, cooldown_millis: u64) -> Self {
|
||||
Self {
|
||||
state: AtomicU8::new(CircuitState::Closed as u8),
|
||||
consecutive_failures: AtomicU32::new(0),
|
||||
last_state_change_millis: AtomicU64::new(0),
|
||||
failure_threshold,
|
||||
cooldown_millis,
|
||||
}
|
||||
}
|
||||
|
||||
pub fn current_state(&self) -> CircuitState {
|
||||
match self.state.load(Ordering::SeqCst) {
|
||||
1 => CircuitState::Open,
|
||||
2 => CircuitState::HalfOpen,
|
||||
_ => CircuitState::Closed,
|
||||
}
|
||||
}
|
||||
|
||||
pub fn can_execute(&self) -> bool {
|
||||
let state = self.current_state();
|
||||
match state {
|
||||
CircuitState::Closed => true,
|
||||
CircuitState::HalfOpen => true,
|
||||
CircuitState::Open => {
|
||||
let now = std::time::SystemTime::now()
|
||||
.duration_since(std::time::UNIX_EPOCH)
|
||||
.unwrap_or_default()
|
||||
.as_millis() as u64;
|
||||
let last = self.last_state_change_millis.load(Ordering::SeqCst);
|
||||
if now.saturating_sub(last) >= self.cooldown_millis {
|
||||
self.state
|
||||
.store(CircuitState::HalfOpen as u8, Ordering::SeqCst);
|
||||
tracing::info!("Neovim RPC circuit breaker transitioned to HalfOpen");
|
||||
true
|
||||
} else {
|
||||
false
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub fn record_success(&self) {
|
||||
self.consecutive_failures.store(0, Ordering::SeqCst);
|
||||
let prev = self.state.swap(CircuitState::Closed as u8, Ordering::SeqCst);
|
||||
if prev != CircuitState::Closed as u8 {
|
||||
tracing::info!("Neovim RPC circuit breaker transitioned to Closed (recovered)");
|
||||
}
|
||||
}
|
||||
|
||||
pub fn record_failure(&self) {
|
||||
let failures = self.consecutive_failures.fetch_add(1, Ordering::SeqCst) + 1;
|
||||
if failures >= self.failure_threshold {
|
||||
let now = std::time::SystemTime::now()
|
||||
.duration_since(std::time::UNIX_EPOCH)
|
||||
.unwrap_or_default()
|
||||
.as_millis() as u64;
|
||||
self.last_state_change_millis.store(now, Ordering::SeqCst);
|
||||
self.state.store(CircuitState::Open as u8, Ordering::SeqCst);
|
||||
tracing::warn!(
|
||||
"Neovim RPC circuit breaker tripped to Open (consecutive failures: {})",
|
||||
failures
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
pub fn reset(&self) {
|
||||
self.consecutive_failures.store(0, Ordering::SeqCst);
|
||||
self.state.store(CircuitState::Closed as u8, Ordering::SeqCst);
|
||||
}
|
||||
}
|
||||
|
||||
pub static CIRCUIT_BREAKER: LazyLock<NvimCircuitBreaker> =
|
||||
LazyLock::new(|| NvimCircuitBreaker::new(2, 5000));
|
||||
|
||||
static NEXT_MSGID: AtomicU64 = AtomicU64::new(1);
|
||||
static RPC_SEMAPHORE: LazyLock<Arc<tokio::sync::Semaphore>> =
|
||||
LazyLock::new(|| Arc::new(tokio::sync::Semaphore::new(100)));
|
||||
@@ -529,6 +620,10 @@ async fn get_nvim_connection() -> Result<mpsc::Sender<NvimRequest>, String> {
|
||||
}
|
||||
|
||||
async fn call_nvim(req: rmpv::Value) -> Result<rmpv::Value, String> {
|
||||
if !CIRCUIT_BREAKER.can_execute() {
|
||||
return Err("Neovim RPC circuit breaker is OPEN (consecutive failures detected). Failing fast.".to_string());
|
||||
}
|
||||
|
||||
let msgid = if let rmpv::Value::Array(ref arr) = req {
|
||||
if arr.len() > 1 {
|
||||
match &arr[1] {
|
||||
@@ -542,7 +637,13 @@ async fn call_nvim(req: rmpv::Value) -> Result<rmpv::Value, String> {
|
||||
0
|
||||
};
|
||||
|
||||
let tx = get_nvim_connection().await?;
|
||||
let tx = match get_nvim_connection().await {
|
||||
Ok(t) => t,
|
||||
Err(e) => {
|
||||
CIRCUIT_BREAKER.record_failure();
|
||||
return Err(e);
|
||||
}
|
||||
};
|
||||
let (reply_tx, reply_rx) = oneshot::channel();
|
||||
|
||||
if let Err(_) = tx
|
||||
@@ -556,20 +657,30 @@ async fn call_nvim(req: rmpv::Value) -> Result<rmpv::Value, String> {
|
||||
PENDING_REQUESTS.remove(&msgid);
|
||||
let mut conn = NVIM_CONN.lock().await;
|
||||
*conn = None;
|
||||
CIRCUIT_BREAKER.record_failure();
|
||||
return Err(
|
||||
"Failed to send request to Neovim connection manager: connection closed".to_string(),
|
||||
);
|
||||
}
|
||||
|
||||
match tokio::time::timeout(tokio::time::Duration::from_secs(15), reply_rx).await {
|
||||
Ok(Ok(res)) => res,
|
||||
match tokio::time::timeout(tokio::time::Duration::from_secs(5), reply_rx).await {
|
||||
Ok(Ok(res)) => {
|
||||
CIRCUIT_BREAKER.record_success();
|
||||
res
|
||||
}
|
||||
Ok(Err(_)) => {
|
||||
PENDING_REQUESTS.remove(&msgid);
|
||||
let mut conn = NVIM_CONN.lock().await;
|
||||
*conn = None;
|
||||
CIRCUIT_BREAKER.record_failure();
|
||||
Err("Response channel dropped".to_string())
|
||||
}
|
||||
Err(_) => {
|
||||
PENDING_REQUESTS.remove(&msgid);
|
||||
Err("Timeout waiting for Neovim response (15s)".to_string())
|
||||
let mut conn = NVIM_CONN.lock().await;
|
||||
*conn = None;
|
||||
CIRCUIT_BREAKER.record_failure();
|
||||
Err("Timeout waiting for Neovim response (5s)".to_string())
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1999,4 +2110,33 @@ mod tests {
|
||||
panic!("Expected array response for nvim_get_api_info");
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_nvim_circuit_breaker_transitions() {
|
||||
let cb = NvimCircuitBreaker::new(2, 50); // 2 failures, 50ms cooldown
|
||||
assert_eq!(cb.current_state(), CircuitState::Closed);
|
||||
assert!(cb.can_execute());
|
||||
|
||||
// 1st failure - remains Closed
|
||||
cb.record_failure();
|
||||
assert_eq!(cb.current_state(), CircuitState::Closed);
|
||||
assert!(cb.can_execute());
|
||||
|
||||
// 2nd failure - trips to Open
|
||||
cb.record_failure();
|
||||
assert_eq!(cb.current_state(), CircuitState::Open);
|
||||
assert!(!cb.can_execute(), "Circuit breaker should fail fast when Open");
|
||||
|
||||
// Wait for cooldown
|
||||
std::thread::sleep(std::time::Duration::from_millis(60));
|
||||
|
||||
// After cooldown, can_execute transitions to HalfOpen
|
||||
assert!(cb.can_execute(), "After cooldown, should allow HalfOpen probe");
|
||||
assert_eq!(cb.current_state(), CircuitState::HalfOpen);
|
||||
|
||||
// Success in HalfOpen recovers back to Closed
|
||||
cb.record_success();
|
||||
assert_eq!(cb.current_state(), CircuitState::Closed);
|
||||
assert!(cb.can_execute());
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user