| 29 | |
| 30 | |
| 31 | class 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) |