(self, api_server_pid=None)
| 1112 | return max(unhandled, 0) |
| 1113 | |
| 1114 | def start_zmq_service(self, api_server_pid=None): |
| 1115 | if api_server_pid is None: |
| 1116 | return |
| 1117 | self.api_server_pid = api_server_pid |
| 1118 | if envs.FD_ENABLE_INTERNAL_ADAPTER: |
| 1119 | self.recv_request_server = ZmqTcpServer(port=envs.FD_ZMQ_RECV_REQUEST_SERVER_PORT, mode=zmq.PULL) |
| 1120 | self.send_response_server = ZmqTcpServer(port=envs.FD_ZMQ_SEND_RESPONSE_SERVER_PORT, mode=zmq.ROUTER) |
| 1121 | self.internal_adapter = InternalAdapter( |
| 1122 | cfg=self.cfg, engine=self, dp_rank=self.cfg.parallel_config.local_data_parallel_id |
| 1123 | ) |
| 1124 | else: |
| 1125 | self.recv_request_server = ZmqIpcServer(name=api_server_pid, mode=zmq.PULL) |
| 1126 | self.send_response_server = ZmqIpcServer(name=api_server_pid, mode=zmq.ROUTER) |
| 1127 | self.recv_result_handle_thread = threading.Thread( |
| 1128 | target=self.send_response_server.recv_result_handle, daemon=True |
| 1129 | ) |
| 1130 | self.recv_result_handle_thread.start() |
| 1131 | time.sleep(3) |
| 1132 | self.insert_task_to_scheduler_thread = threading.Thread(target=self._insert_zmq_task_to_scheduler, daemon=True) |
| 1133 | self.insert_task_to_scheduler_thread.start() |
| 1134 | |
| 1135 | self.receive_output_thread = threading.Thread(target=self._zmq_send_generated_tokens, daemon=True) |
| 1136 | self.receive_output_thread.start() |
| 1137 | |
| 1138 | def _insert_zmq_task_to_scheduler(self): |
| 1139 | tracing.trace_set_thread_info("Insert Task to Scheduler") |
no test coverage detected