(mut self)
| 139 | } |
| 140 | |
| 141 | async fn run(mut self) { |
| 142 | loop { |
| 143 | let res = if let Some(current) = &mut self.current_state { |
| 144 | tokio::select! { |
| 145 | scheduler_cmd = self.scheduler_rx.recv() => { |
| 146 | if let Some(cmd) = scheduler_cmd{ |
| 147 | self.handle_scheduler_cmd(cmd).await |
| 148 | }else{ |
| 149 | Ok(ContinueState::Stop) |
| 150 | } |
| 151 | } |
| 152 | runtime_event = self.runtime_evt_rx.recv() => { |
| 153 | if let Some(evt) = runtime_event{ |
| 154 | self.handle_runtime_evt(evt).await |
| 155 | }else{ |
| 156 | Ok(ContinueState::Stop) |
| 157 | } |
| 158 | } |
| 159 | vm_event = current.evt_rx.recv() => { |
| 160 | if let Some(evt) = vm_event{ |
| 161 | self.handle_vm_evt(evt).await |
| 162 | }else{ |
| 163 | info!("vm shut down: channel closed"); |
| 164 | let current = self.current_state.take().unwrap(); |
| 165 | self.handle_vm_channel_closed(¤t).await.map(|_|ContinueState::Continue) |
| 166 | } |
| 167 | } |
| 168 | } |
| 169 | } else { |
| 170 | tokio::select! { |
| 171 | scheduler_cmd = self.scheduler_rx.recv() => { |
| 172 | if let Some(cmd) = scheduler_cmd{ |
| 173 | self.handle_scheduler_cmd(cmd).await |
| 174 | }else{ |
| 175 | Ok(ContinueState::Stop) |
| 176 | } |
| 177 | } |
| 178 | runtime_event = self.runtime_evt_rx.recv() => { |
| 179 | if let Some(evt) = runtime_event{ |
| 180 | self.handle_runtime_evt(evt).await |
| 181 | }else{ |
| 182 | Ok(ContinueState::Stop) |
| 183 | } |
| 184 | } |
| 185 | } |
| 186 | }; |
| 187 | |
| 188 | match res { |
| 189 | Err(err) => { |
| 190 | error!(%err, "failed sending scheduler message") |
| 191 | } |
| 192 | Ok(ContinueState::Stop) => break, |
| 193 | Ok(ContinueState::Continue) => {} |
| 194 | } |
| 195 | } |
| 196 | |
| 197 | if let Err(err) = self.wait_shutdown_current_vm().await { |
| 198 | error!(%err, "failed shutting down current vm") |
no test coverage detected