Allocate resource and insert tasks to engine. Used in v0_kvcache_scheduler.
(self, tasks: List[Request], current_id=-1)
| 460 | ) |
| 461 | |
| 462 | def insert_tasks(self, tasks: List[Request], current_id=-1): |
| 463 | """ |
| 464 | Allocate resource and insert tasks to engine. |
| 465 | Used in v0_kvcache_scheduler. |
| 466 | """ |
| 467 | if not isinstance(tasks, list): |
| 468 | tasks = [tasks] |
| 469 | |
| 470 | self.resource_manager.check_and_free_block_tables() |
| 471 | |
| 472 | need_delete_tasks = [] |
| 473 | for task in tasks: |
| 474 | rid = task.request_id.split("_")[0] |
| 475 | trace_carrier = task.trace_carrier |
| 476 | if trace_carrier: |
| 477 | tracing.trace_set_proc_propagate_context(rid, trace_carrier) |
| 478 | task.trace_carrier = tracing.trace_get_proc_propagate_context(rid) |
| 479 | if self.cfg.scheduler_config.splitwise_role == "prefill": |
| 480 | status, msg = self.split_connector.check_decode_allocated(task) |
| 481 | if status: |
| 482 | task.metrics.ask_decode_resource_finish_time = time.time() |
| 483 | else: |
| 484 | self.llm_logger.error(f"{task.request_id} prefill failed with msg:{msg}.") |
| 485 | self.scheduler.put_results( |
| 486 | [ |
| 487 | RequestOutput( |
| 488 | request_id=task.request_id, |
| 489 | finished=True, |
| 490 | error_code=500, |
| 491 | error_msg=msg, |
| 492 | ) |
| 493 | ] |
| 494 | ) |
| 495 | need_delete_tasks.append(task) |
| 496 | continue |
| 497 | for tmp_task in need_delete_tasks: |
| 498 | tasks.remove(tmp_task) |
| 499 | |
| 500 | for item in tasks: |
| 501 | trace_print(LoggingEventName.RESOURCE_ALLOCATE_START, item.request_id, getattr(item, "user", "")) |
| 502 | |
| 503 | available_batch = np.sum(self.resource_manager.stop_flags) |
| 504 | if len(tasks) > available_batch: |
| 505 | self.llm_logger.error(f"Inserting batch:{len(tasks)} exceeds the available batch:{available_batch}.") |
| 506 | self.llm_logger.error("The exceeded part will be ignored!") |
| 507 | tasks = tasks[:available_batch] |
| 508 | |
| 509 | req_ids = [t.request_id for t in tasks] |
| 510 | |
| 511 | tasks = self.resource_manager.allocate_resources_for_new_tasks(tasks) |
| 512 | |
| 513 | if not tasks: |
| 514 | error_msg = f"The request required resources is exceed the limit, request id={req_ids}." |
| 515 | self.llm_logger.error(error_msg) |
| 516 | raise EngineError(error_msg, error_code=500) |
| 517 | return False |
| 518 | |
| 519 | self.token_processor.number_of_tasks += len(tasks) |
no test coverage detected