Control function for resuming request generation. This method resumes the paused request generation process by setting the pause flag and notifying all waiting threads. It logs the start and end of the resume operation. Args: control_request: Control request obj
(self, control_request: ControlRequest)
| 1354 | return None |
| 1355 | |
| 1356 | def _control_resume(self, control_request: ControlRequest) -> Optional[dict]: |
| 1357 | """Control function for resuming request generation. |
| 1358 | |
| 1359 | This method resumes the paused request generation process by setting the pause flag |
| 1360 | and notifying all waiting threads. It logs the start and end of the resume operation. |
| 1361 | |
| 1362 | Args: |
| 1363 | control_request: Control request object containing resume operation information |
| 1364 | """ |
| 1365 | self.llm_logger.info("Start to resume request generation.") |
| 1366 | with self._pause_cond: |
| 1367 | if not self.is_paused: |
| 1368 | self.llm_logger.info("Engine is not paused, no need to resume.") |
| 1369 | return None |
| 1370 | self.is_paused = False |
| 1371 | self._pause_cond.notify_all() |
| 1372 | |
| 1373 | # resume cache transfer |
| 1374 | if self.cfg.cache_config.num_cpu_blocks > 0 or self.cfg.cache_config.kvcache_storage_backend: |
| 1375 | self.llm_logger.info("Start to resume cache transfer.") |
| 1376 | resume_transfer_request = ControlRequest(request_id="resume_transfer", method="resume") |
| 1377 | self.cache_task_queue.put_transfer_task((CacheStatus.CTRL, resume_transfer_request)) |
| 1378 | # Wait for cache_transfer responses |
| 1379 | asyncio.run(self._wait_for_control_responses("resume_transfer", 60, executors=["cache_transfer"])) |
| 1380 | self.llm_logger.info("Successfully resumed cache transfer.") |
| 1381 | |
| 1382 | self.llm_logger.info("Successfully resumed request generation.") |
| 1383 | return None |
| 1384 | |
| 1385 | def _control_is_paused(self, control_request: ControlRequest) -> bool: |
| 1386 | """ |
no test coverage detected