MCPcopy Create free account
hub / github.com/PaddlePaddle/FastDeploy / insert_tasks

Method insert_tasks

fastdeploy/engine/common_engine.py:462–561  ·  view source on GitHub ↗

Allocate resource and insert tasks to engine. Used in v0_kvcache_scheduler.

(self, tasks: List[Request], current_id=-1)

Source from the content-addressed store, hash-verified

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)

Calls 14

RequestOutputClass · 0.90
EngineErrorClass · 0.90
splitMethod · 0.80
put_tasksMethod · 0.80
errorMethod · 0.45

Tested by

no test coverage detected