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

Class InternalAdapter

fastdeploy/splitwise/internal_adapter_utils.py:31–121  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

29
30
31class InternalAdapter:
32 def __init__(self, cfg, engine, dp_rank):
33 self.cfg = cfg
34 self.engine = engine
35 self.dp_rank = dp_rank
36 recv_control_cmd_ports = envs.FD_ZMQ_CONTROL_CMD_SERVER_PORTS.split(",")
37 self.response_lock = threading.Lock() # prevent to call send_multipart in zmq concurrently
38 self.recv_control_cmd_server = ZmqTcpServer(port=recv_control_cmd_ports[dp_rank], mode=zmq.ROUTER)
39 self.recv_external_instruct_thread = threading.Thread(
40 target=self._recv_external_module_control_instruct, daemon=True
41 )
42 self.recv_external_instruct_thread.start()
43 if cfg.scheduler_config.splitwise_role != "mixed":
44 self.response_external_instruct_thread = threading.Thread(
45 target=self._response_external_module_control_instruct, daemon=True
46 )
47 self.response_external_instruct_thread.start()
48
49 def _get_current_server_info(self):
50 """
51 Get resources information
52 """
53 available_batch_size = min(self.cfg.max_prefill_batch, self.engine.resource_manager.available_batch())
54
55 available_block_num = self.engine.resource_manager.available_block_num()
56 server_info = {
57 "splitwise_role": self.cfg.scheduler_config.splitwise_role,
58 "block_size": int(self.cfg.cache_config.block_size),
59 "block_num": int(available_block_num),
60 "max_block_num": int(self.cfg.cache_config.total_block_num),
61 "dec_token_num": int(self.cfg.cache_config.dec_token_num),
62 "available_resource": float(1.0 * available_block_num / self.cfg.cache_config.total_block_num),
63 "max_batch_size": int(available_batch_size),
64 "max_input_token_num": self.cfg.model_config.max_model_len,
65 "unhandled_request_num": self.engine.scheduler.get_unhandled_request_num(),
66 "available_batch": int(self.engine.resource_manager.available_batch()),
67 }
68 return server_info
69
70 def _recv_external_module_control_instruct(self):
71 """
72 Receive a multipart message from the control cmd socket.
73 """
74 while True:
75 try:
76 with self.response_lock:
77 task = self.recv_control_cmd_server.recv_control_cmd()
78 if task is None:
79 time.sleep(0.001)
80 continue
81 logger.info(f"dprank {self.dp_rank} Recieve control task: {task}")
82 task_id_str = task["task_id"]
83 if task["cmd"] == "get_payload":
84 payload_info = self._get_current_server_info()
85 result = {"task_id": task_id_str, "result": payload_info}
86 logger.debug(f"Response for task: {task_id_str}")
87 with self.response_lock:
88 self.recv_control_cmd_server.response_for_control_cmd(task_id_str, result)

Callers 1

start_zmq_serviceMethod · 0.90

Calls

no outgoing calls

Tested by

no test coverage detected