(&mut self, ev: NetworkEvent)
| 583 | |
| 584 | #[tracing::instrument(skip(self))] |
| 585 | fn emit_network_event(&mut self, ev: NetworkEvent) { |
| 586 | let mut to_remove = Vec::new(); |
| 587 | for (i, (open, sender)) in self.network_events.iter_mut().enumerate() { |
| 588 | if !open.load(Ordering::Relaxed) { |
| 589 | to_remove.push(i); |
| 590 | continue; |
| 591 | } |
| 592 | let ev = ev.clone(); |
| 593 | let sender = sender.clone(); |
| 594 | let open = open.clone(); |
| 595 | tokio::task::spawn(async move { |
| 596 | if let Err(_e) = sender.send(ev.clone()).await { |
| 597 | // Mark sender as closed so we stop sending events to it |
| 598 | open.store(false, Ordering::Relaxed); |
| 599 | } |
| 600 | }); |
| 601 | } |
| 602 | for idx in to_remove.iter().rev() { |
| 603 | self.network_events.swap_remove(*idx); |
| 604 | } |
| 605 | } |
| 606 | |
| 607 | #[tracing::instrument(skip_all)] |
| 608 | async fn handle_node_event(&mut self, event: Event) -> Result<Option<SwarmEventResult>> { |
no test coverage detected