(self)
| 123 | atexit.register(self.pre_shutdown) |
| 124 | |
| 125 | def _setup_queues(self) -> WorkerCommIpcAddrs: |
| 126 | |
| 127 | self.request_queue = IpcQueue(is_server=True, |
| 128 | name="proxy_request_queue") |
| 129 | self.worker_init_status_queue = IpcQueue( |
| 130 | is_server=True, |
| 131 | socket_type=zmq.ROUTER, |
| 132 | name="worker_init_status_queue") |
| 133 | # TODO[chunweiy]: Unify IpcQueue and FusedIpcQueue |
| 134 | # Use PULL mode when enable_postprocess_parallel as there are |
| 135 | # multiple senders from multiple processes. |
| 136 | self.result_queue = FusedIpcQueue( |
| 137 | is_server=True, |
| 138 | fuse_message=False, |
| 139 | socket_type=zmq.PULL |
| 140 | if self.enable_postprocess_parallel else zmq.PAIR, |
| 141 | name="proxy_result_queue") |
| 142 | # Stats and KV events are now fetched via RPC, not IPC queues. |
| 143 | return WorkerCommIpcAddrs( |
| 144 | request_queue_addr=self.request_queue.address, |
| 145 | worker_init_status_queue_addr=self.worker_init_status_queue.address, |
| 146 | result_queue_addr=self.result_queue.address, |
| 147 | ) |
| 148 | |
| 149 | def abort_request(self, request_id: int) -> None: |
| 150 | ''' Abort a request by sending a cancelling request to the request queue. |
no test coverage detected