Start a worker thread to broadcast an event to each RateLimiterGroupHandle when the RateLimiter becomes unblocked.
(&mut self, exit_evt: EventFd)
| 217 | /// Start a worker thread to broadcast an event to each RateLimiterGroupHandle |
| 218 | /// when the RateLimiter becomes unblocked. |
| 219 | pub fn start_thread(&mut self, exit_evt: EventFd) -> result::Result<(), Error> { |
| 220 | let inner = self.inner.clone(); |
| 221 | let epoll_fd = self.epoll_file.as_raw_fd(); |
| 222 | thread::Builder::new() |
| 223 | .name(format!("rate-limit-group-{}", inner.id)) |
| 224 | .spawn(move || { |
| 225 | let res = std::panic::catch_unwind(AssertUnwindSafe(move || { |
| 226 | const EPOLL_EVENTS_LEN: usize = 2; |
| 227 | |
| 228 | let mut events = |
| 229 | [epoll::Event::new(epoll::Events::empty(), 0); EPOLL_EVENTS_LEN]; |
| 230 | |
| 231 | loop { |
| 232 | let num_events = match epoll::wait(epoll_fd, -1, &mut events[..]) { |
| 233 | Ok(res) => res, |
| 234 | Err(e) => { |
| 235 | if e.kind() == io::ErrorKind::Interrupted { |
| 236 | continue; |
| 237 | } |
| 238 | return Err(Error::Epoll(e)); |
| 239 | } |
| 240 | }; |
| 241 | |
| 242 | for event in events.iter().take(num_events) { |
| 243 | let dispatch_event: EpollDispatch = event.data.into(); |
| 244 | match dispatch_event { |
| 245 | EpollDispatch::Unknown => { |
| 246 | let event = event.data; |
| 247 | warn!("Unknown rate-limiter loop event: {event}"); |
| 248 | } |
| 249 | EpollDispatch::Unblocked => { |
| 250 | inner.rate_limiter.event_handler().unwrap(); |
| 251 | let handles = inner.handles.lock().unwrap(); |
| 252 | for handle in handles.iter() { |
| 253 | handle.write(1).map_err(Error::EventFdWrite)?; |
| 254 | } |
| 255 | } |
| 256 | EpollDispatch::Kill => { |
| 257 | info!( |
| 258 | "KILL_EVENT received, stopping rate-limit-group epoll loop" |
| 259 | ); |
| 260 | return Ok(()); |
| 261 | } |
| 262 | } |
| 263 | } |
| 264 | } |
| 265 | })); |
| 266 | |
| 267 | match res { |
| 268 | Ok(res) => { |
| 269 | if let Err(e) = res { |
| 270 | error!("Error running rate-limit-group worker: {e:?}"); |
| 271 | exit_evt.write(1).unwrap(); |
| 272 | } |
| 273 | } |
| 274 | Err(_) => { |
| 275 | error!("rate-limit-group worker panicked"); |
| 276 | exit_evt.write(1).unwrap(); |