(mut self)
| 41 | |
| 42 | impl BrokerConn { |
| 43 | async fn run(mut self) -> std::io::Result<()> { |
| 44 | let _ = self.scheduler_tx.send(SchedulerCommand::BrokerConnected); |
| 45 | |
| 46 | loop { |
| 47 | let next: BrokerEvent = simpleproto::read_message(&mut self.stream).await?; |
| 48 | match self.handle_broker_message(next).await? { |
| 49 | ContinueState::Continue => {} |
| 50 | ContinueState::Stop => return Ok(()), |
| 51 | } |
| 52 | } |
| 53 | } |
| 54 | |
| 55 | #[instrument(skip(self, message))] |
| 56 | async fn handle_broker_message( |
no test coverage detected