Skip to content
Merged
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
71 changes: 59 additions & 12 deletions src/discord.rs
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
use std::fmt;

use std::collections::HashMap;
use std::sync::{Arc, Mutex};
use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
use std::time::Duration;

use reqwest::header::{AUTHORIZATION, CONTENT_TYPE, HeaderMap, HeaderValue};
Expand Down Expand Up @@ -737,14 +737,27 @@ impl DiscordClient {
})
}

/// Borrow the shared Discord delivery state, tolerating lock poisoning.
///
/// A panic inside one Discord critical section must not permanently brick
/// delivery for the rest of the daemon's lifetime. The guarded state is a
/// rate limiter, per-target circuit breakers, and the DLQ buffer: all
/// individually recoverable, so recovering the poisoned guard degrades to
/// possibly-stale counters instead of an unrecoverable panic loop on every
/// later send. This matches the poison-tolerant locking already used by
/// the daemon, dispatch, lane, and subscription paths.
fn state(&self) -> MutexGuard<'_, DiscordState> {
self.state.lock().unwrap_or_else(PoisonError::into_inner)
}

fn allow_request(
&self,
key: &str,
) -> (
bool,
Option<crate::core::circuit_breaker::CircuitTransition>,
) {
let mut state = self.state.lock().expect("discord state lock");
let mut state = self.state();
state
.circuits
.entry(key.to_string())
Expand All @@ -758,12 +771,12 @@ impl DiscordClient {
}

fn rate_limit_delay(&self, key: &str) -> Duration {
let mut state = self.state.lock().expect("discord state lock");
let mut state = self.state();
state.limiter.delay_for(key)
}

fn record_success(&self, key: &str) -> Option<crate::core::circuit_breaker::CircuitTransition> {
let mut state = self.state.lock().expect("discord state lock");
let mut state = self.state();
state
.circuits
.entry(key.to_string())
Expand All @@ -777,7 +790,7 @@ impl DiscordClient {
}

fn record_failure(&self, key: &str) -> Option<crate::core::circuit_breaker::CircuitTransition> {
let mut state = self.state.lock().expect("discord state lock");
let mut state = self.state();
state
.circuits
.entry(key.to_string())
Expand Down Expand Up @@ -842,7 +855,7 @@ impl DiscordClient {
.unwrap_or_else(|_| "{\"error\":\"dlq serialize failed\"}".to_string())
);

let mut state = self.state.lock().expect("discord state lock");
let mut state = self.state();
state.dlq.push(entry);
}

Expand Down Expand Up @@ -876,12 +889,17 @@ impl DiscordClient {

#[cfg(test)]
fn dlq_entries(&self) -> Vec<DlqEntry> {
self.state
.lock()
.expect("discord state lock")
.dlq
.entries()
.to_vec()
self.state().dlq.entries().to_vec()
}

#[cfg(test)]
fn poison_state_for_tests(&self) {
let state = Arc::clone(&self.state);
let _ = std::thread::spawn(move || {
let _guard = state.lock().unwrap_or_else(PoisonError::into_inner);
panic!("poison discord state");
})
.join();
}
}

Expand Down Expand Up @@ -1581,6 +1599,35 @@ mod tests {
assert!(!format!("{error:?}").contains(sentinel));
}

#[test]
fn poisoned_delivery_state_still_serves_limiter_circuit_and_dlq() {
let client =
DiscordClient::for_tests_with_api_base("test-token", "http://127.0.0.1:1".to_string())
.unwrap();

// Baseline: a healthy client admits requests and holds an empty DLQ.
let (allowed, _) = client.allow_request("channel:poison");
assert!(allowed);

client.poison_state_for_tests();

// After poisoning, every state-touching path must keep working instead
// of panicking for the rest of the daemon's lifetime.
let (allowed_after, _) = client.allow_request("channel:poison");
assert!(allowed_after);
assert_eq!(client.rate_limit_delay("channel:poison"), Duration::ZERO);
assert!(client.record_success("channel:poison").is_none());
for _ in 0..CIRCUIT_FAILURE_THRESHOLD {
client.record_failure("channel:poison");
}
let (allowed_open, _) = client.allow_request("channel:poison");
assert!(
!allowed_open,
"circuit must still open on a recovered poisoned state"
);
assert!(client.dlq_entries().is_empty());
}

#[tokio::test]
async fn lane_read_probes_are_content_free_and_do_not_touch_delivery_state() {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
Expand Down
Loading