running_process/broker/server/
fd_pressure.rs1use std::io;
16use std::sync::Mutex;
17
18use crate::broker::protocol::{ErrorCode, HelloReply};
19
20use super::connection::refused_reply;
21
22pub const DEFAULT_FD_PRESSURE_RECOVERY_ACCEPTS: u32 = 3;
24pub const DEFAULT_FD_PRESSURE_RETRY_AFTER_MS: u64 = 1_000;
26
27#[derive(Clone, Copy, Debug, PartialEq, Eq)]
29pub enum FdPressureDecision {
30 Demoted,
32 Unrelated,
34}
35
36#[derive(Clone, Copy, Debug)]
38pub struct FdPressureConfig {
39 pub recovery_accepts: u32,
41 pub retry_after_ms: u64,
43}
44
45impl Default for FdPressureConfig {
46 fn default() -> Self {
47 Self {
48 recovery_accepts: DEFAULT_FD_PRESSURE_RECOVERY_ACCEPTS,
49 retry_after_ms: DEFAULT_FD_PRESSURE_RETRY_AFTER_MS,
50 }
51 }
52}
53
54#[derive(Clone, Copy, Debug, Default)]
55struct GuardState {
56 demoted: bool,
57 consecutive_ok: u32,
58 demotions_total: u64,
59 refused_while_demoted: u64,
60}
61
62#[derive(Debug, Default)]
64pub struct FdPressureGuard {
65 config: FdPressureConfig,
66 state: Mutex<GuardState>,
67}
68
69impl FdPressureGuard {
70 pub fn new(config: FdPressureConfig) -> Self {
72 Self {
73 config,
74 state: Mutex::new(GuardState::default()),
75 }
76 }
77
78 pub fn on_accept_error(&self, err: &io::Error) -> FdPressureDecision {
81 if !is_fd_exhaustion_error(err) {
82 return FdPressureDecision::Unrelated;
83 }
84 let mut state = self.lock();
85 if !state.demoted {
86 state.demoted = true;
87 state.demotions_total += 1;
88 }
89 state.consecutive_ok = 0;
90 FdPressureDecision::Demoted
91 }
92
93 pub fn on_accept_ok(&self) -> bool {
96 let mut state = self.lock();
97 if !state.demoted {
98 return false;
99 }
100 state.consecutive_ok += 1;
101 if state.consecutive_ok >= self.config.recovery_accepts {
102 state.demoted = false;
103 state.consecutive_ok = 0;
104 return true;
105 }
106 false
107 }
108
109 pub fn is_demoted(&self) -> bool {
111 self.lock().demoted
112 }
113
114 pub fn demotions_total(&self) -> u64 {
116 self.lock().demotions_total
117 }
118
119 pub fn refused_while_demoted(&self) -> u64 {
121 self.lock().refused_while_demoted
122 }
123
124 pub fn refusal_reply(&self) -> HelloReply {
126 self.lock().refused_while_demoted += 1;
127 refused_reply(
128 ErrorCode::ErrorFdPressure,
129 "broker is low on file descriptors; retry shortly",
130 self.config.retry_after_ms,
131 )
132 }
133
134 pub fn force_demote(&self) {
136 let mut state = self.lock();
137 if !state.demoted {
138 state.demoted = true;
139 state.demotions_total += 1;
140 }
141 state.consecutive_ok = 0;
142 }
143
144 fn lock(&self) -> std::sync::MutexGuard<'_, GuardState> {
145 self.state
146 .lock()
147 .unwrap_or_else(|poisoned| poisoned.into_inner())
148 }
149}
150
151pub fn is_fd_exhaustion_error(err: &io::Error) -> bool {
153 crate::platform::resources::signals_fd_exhaustion(err)
154}
155
156pub fn fd_exhaustion_error_for_tests() -> io::Error {
158 crate::platform::resources::fd_exhaustion_error()
159}
160
161#[cfg(test)]
162mod tests {
163 use super::*;
164 use crate::broker::protocol::hello_reply::Result as HelloReplyResult;
165
166 #[test]
167 fn unrelated_errors_do_not_demote() {
168 let guard = FdPressureGuard::default();
169 let err = io::Error::new(io::ErrorKind::PermissionDenied, "denied");
170 assert_eq!(guard.on_accept_error(&err), FdPressureDecision::Unrelated);
171 assert!(!guard.is_demoted());
172 assert_eq!(guard.demotions_total(), 0);
173 }
174
175 #[test]
176 fn fd_exhaustion_demotes_and_recovers_after_streak() {
177 let guard = FdPressureGuard::new(FdPressureConfig {
178 recovery_accepts: 2,
179 retry_after_ms: 250,
180 });
181 assert_eq!(
182 guard.on_accept_error(&fd_exhaustion_error_for_tests()),
183 FdPressureDecision::Demoted
184 );
185 assert!(guard.is_demoted());
186 assert_eq!(guard.demotions_total(), 1);
187
188 assert!(!guard.on_accept_ok());
189 assert!(guard.is_demoted());
190 assert!(guard.on_accept_ok());
191 assert!(!guard.is_demoted());
192 }
193
194 #[test]
195 fn accept_error_resets_recovery_streak() {
196 let guard = FdPressureGuard::new(FdPressureConfig {
197 recovery_accepts: 2,
198 retry_after_ms: 250,
199 });
200 guard.on_accept_error(&fd_exhaustion_error_for_tests());
201 assert!(!guard.on_accept_ok());
202 guard.on_accept_error(&fd_exhaustion_error_for_tests());
203 assert!(!guard.on_accept_ok());
204 assert!(guard.is_demoted());
205 assert!(guard.on_accept_ok());
206 assert!(!guard.is_demoted());
207 assert_eq!(guard.demotions_total(), 1);
208 }
209
210 #[test]
211 fn refusal_reply_uses_reserved_fd_pressure_code() {
212 let guard = FdPressureGuard::default();
213 guard.force_demote();
214 let reply = guard.refusal_reply();
215 let HelloReplyResult::Refused(refused) = reply.result.unwrap() else {
216 panic!("expected refusal");
217 };
218 assert_eq!(
219 ErrorCode::try_from(refused.code),
220 Ok(ErrorCode::ErrorFdPressure)
221 );
222 assert_eq!(refused.retry_after_ms, DEFAULT_FD_PRESSURE_RETRY_AFTER_MS);
223 assert_eq!(guard.refused_while_demoted(), 1);
224 }
225
226 #[test]
227 fn ok_accepts_while_healthy_are_no_ops() {
228 let guard = FdPressureGuard::default();
229 assert!(!guard.on_accept_ok());
230 assert!(!guard.is_demoted());
231 }
232}