(
&self,
payload: EventPayload,
now: Option<DateTime<Utc>>,
forward_mode: ForwardMode,
)
| 362 | } |
| 363 | |
| 364 | fn send_event( |
| 365 | &self, |
| 366 | payload: EventPayload, |
| 367 | now: Option<DateTime<Utc>>, |
| 368 | forward_mode: ForwardMode, |
| 369 | ) { |
| 370 | let mut inner = self.inner.get_mut(); |
| 371 | if inner.storage_sender.is_none() && inner.client_listeners.is_empty() { |
| 372 | return; |
| 373 | } |
| 374 | let event = Event { |
| 375 | time: now.unwrap_or_else(Utc::now), |
| 376 | payload, |
| 377 | }; |
| 378 | |
| 379 | inner.client_listeners.retain(|listener| { |
| 380 | if listener.filter.check(&event.payload) { |
| 381 | listener.sender.send(event.clone()).is_ok() |
| 382 | } else { |
| 383 | true |
| 384 | } |
| 385 | }); |
| 386 | |
| 387 | if let Some(ref streamer) = inner.storage_sender |
| 388 | && matches!(forward_mode, ForwardMode::StreamAndPersist) |
| 389 | && streamer.send(EventStreamMessage::Event(event)).is_err() |
| 390 | { |
| 391 | log::error!("Event streaming queue has been closed."); |
| 392 | } |
| 393 | } |
| 394 | |
| 395 | pub fn on_server_stop(&self) { |
| 396 | self.send_event( |
no test coverage detected