(&mut self)
| 445 | /// workers until one accepts. |
| 446 | #[instrument(skip_all)] |
| 447 | async fn ensure_session(&mut self) { |
| 448 | if self.session.is_some() { |
| 449 | return; |
| 450 | } |
| 451 | |
| 452 | loop { |
| 453 | let worker = self |
| 454 | .worker_pool |
| 455 | .req_worker(self.guild_id, self.premium_tier.option()) |
| 456 | .await; |
| 457 | |
| 458 | info!(tier = worker.priority_index, "claimed new worker"); |
| 459 | self.last_claimed_worker_at = Instant::now(); |
| 460 | |
| 461 | let can_resume = !self.config.no_reuse_vms |
| 462 | && !self.force_load_scripts_next |
| 463 | && matches!(&worker.session_state, Some(s) if s.guild_id == self.guild_id); |
| 464 | |
| 465 | if can_resume { |
| 466 | self.session = Some(VmSession::resume( |
| 467 | worker, |
| 468 | self.guild_id, |
| 469 | self.logger.clone(), |
| 470 | )); |
| 471 | } else { |
| 472 | // a new vm will send us a new set of timers and tasks |
| 473 | self.clear_loaded_timers_and_tasks(); |
| 474 | |
| 475 | self.vm_session_id_gen += 1; |
| 476 | match VmSession::create( |
| 477 | worker, |
| 478 | self.guild_id, |
| 479 | self.logger.clone(), |
| 480 | self.vm_session_id_gen, |
| 481 | self.scripts.clone(), |
| 482 | self.premium_tier.option(), |
| 483 | ) { |
| 484 | Ok(session) => self.session = Some(session), |
| 485 | Err(worker) => { |
| 486 | self.note_worker_returned(worker.worker_id); |
| 487 | self.worker_pool.return_worker(worker, true); |
| 488 | continue; |
| 489 | } |
| 490 | } |
| 491 | } |
| 492 | |
| 493 | self.force_load_scripts_next = false; |
| 494 | return; |
| 495 | } |
| 496 | } |
| 497 | |
| 498 | async fn discard_broken_session(&mut self) { |
| 499 | // reason-based scheduler notifications are still sent from within, |
no test coverage detected