MCPcopy Create free account
hub / github.com/Botloader/botloader / run

Method run

components/vmworker/src/lib.rs:141–200  ·  view source on GitHub ↗
(mut self)

Source from the content-addressed store, hash-verified

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(&current).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")

Callers 1

runFunction · 0.45

Calls 1

Tested by

no test coverage detected