We have a new deployment to execute, but no task yet. Contact the gateway to claim the deployment and get a task to execute
(
&mut self,
mut deployment: DeploymentContext,
)
| 271 | |
| 272 | /// We have a new deployment to execute, but no task yet. Contact the gateway to claim the deployment and get a task to execute |
| 273 | async fn execute_deployment( |
| 274 | &mut self, |
| 275 | mut deployment: DeploymentContext, |
| 276 | ) -> (DeploymentManagerState, Option<Duration>) { |
| 277 | info!("Starting to execute deployment"); |
| 278 | let (mut msg_stream, close_upstream_tx, logger, metrics_registry, log_file_writer) = |
| 279 | match deployment.execute_deployment(&mut self.engine_client).await { |
| 280 | Ok(upstream_msg) => upstream_msg, |
| 281 | Err(err) => { |
| 282 | return match err.code() { |
| 283 | Code::NotFound | Code::DeadlineExceeded => { |
| 284 | error!( |
| 285 | "Deployment cannot be executed due to {:?} for {:?}", |
| 286 | err, &deployment.deployment_info |
| 287 | ); |
| 288 | (DeploymentManagerState::SeekingNewDeployment {}, None) |
| 289 | } |
| 290 | _ => { |
| 291 | error!("Error while getting new deployment: {}", err); |
| 292 | let next_step = DeploymentManagerState::ExecutingDeployment { deployment }; |
| 293 | (next_step, Some(self.default_wait_time)) |
| 294 | } |
| 295 | }; |
| 296 | } |
| 297 | }; |
| 298 | info!( |
| 299 | "Connected to gateway, executing deployment task for: {:?}", |
| 300 | deployment.deployment_info |
| 301 | ); |
| 302 | |
| 303 | tokio::select! { |
| 304 | biased; |
| 305 | |
| 306 | // We lost the connection with gateway to forward engine message, trying to reconnect |
| 307 | _ = close_upstream_tx.closed() => { |
| 308 | info!("EngineEvent forwarder to gateway has been close, trying to resume connection"); |
| 309 | let next_step = DeploymentManagerState::ExecutingDeployment { deployment }; |
| 310 | (next_step, None) |
| 311 | } |
| 312 | msg = msg_stream.next() => match msg { |
| 313 | Some(Ok(msg)) => { |
| 314 | // We record the last message we received, so in case of cnx loss |
| 315 | // we can resume the deployment and restart from the last message |
| 316 | deployment.set_last_message_id(msg.message_id.clone()); |
| 317 | |
| 318 | match msg.request { |
| 319 | Some(engine_message_rx::Request::DeploymentRequest(payload)) => { |
| 320 | info!("Received new deployment task: {}", payload); |
| 321 | let task = (self.mk_engine_task)(payload, &deployment.deployment_info, &self.engine_client, logger.clone(), metrics_registry.clone(), log_file_writer.clone()); |
| 322 | let upstream = UpstreamGatewayContext::new(msg_stream, close_upstream_tx, logger, metrics_registry, log_file_writer); |
| 323 | match task { |
| 324 | Ok(task) => { |
| 325 | let next_step = DeploymentManagerState::ExecutingDeploymentTask { |
| 326 | deployment, |
| 327 | task: TaskContext::spawn_new_task(task), |
| 328 | upstream_gtw: Box::new(upstream), |
| 329 | }; |
| 330 | (next_step, None) |
no outgoing calls
no test coverage detected