ZmqTcpServer, used when FD_ENABLE_INTERNAL_ADAPTER=1
| 393 | |
| 394 | |
| 395 | class ZmqTcpServer(ZmqServerBase): |
| 396 | """ |
| 397 | ZmqTcpServer, used when FD_ENABLE_INTERNAL_ADAPTER=1 |
| 398 | """ |
| 399 | |
| 400 | def __init__(self, port, mode): |
| 401 | super(ZmqTcpServer, self).__init__() |
| 402 | self.mode = mode |
| 403 | self.port = port |
| 404 | self.cached_results = defaultdict(list) |
| 405 | self.ZMQ_SNDHWM = int(envs.FD_ZMQ_SNDHWM) |
| 406 | self.aggregate_send = envs.FD_USE_AGGREGATE_SEND |
| 407 | |
| 408 | self.mutex = threading.Lock() |
| 409 | self.req_dict = dict() |
| 410 | self.running = True |
| 411 | self.context = zmq.Context() |
| 412 | self._create_socket() |
| 413 | self.response_token_lock = threading.Lock() |
| 414 | |
| 415 | def _create_socket(self): |
| 416 | """create and return a ZeroMQ socket.""" |
| 417 | self.socket = self.context.socket(self.mode) |
| 418 | self.socket.setsockopt(zmq.SNDHWM, self.ZMQ_SNDHWM) |
| 419 | self.socket.setsockopt(zmq.SNDTIMEO, -1) |
| 420 | self.address = f"tcp://*:{self.port}" |
| 421 | self.socket.bind(self.address) |
| 422 | return self.socket |
| 423 | |
| 424 | def recv_control_cmd(self): |
| 425 | """ |
| 426 | Recieve control command from client |
| 427 | """ |
| 428 | self._ensure_socket() |
| 429 | try: |
| 430 | client, _, task_data = self.socket.recv_multipart(flags=zmq.NOBLOCK) |
| 431 | task = msgpack.unpackb(task_data) |
| 432 | task_id_str = task["task_id"] |
| 433 | except zmq.Again: |
| 434 | return None |
| 435 | with self.mutex: |
| 436 | self.req_dict[task_id_str] = client |
| 437 | return task |
| 438 | |
| 439 | def response_for_control_cmd(self, task_id, result): |
| 440 | """ |
| 441 | Send command result back to client. |
| 442 | """ |
| 443 | self._ensure_socket() |
| 444 | if self.socket is None: |
| 445 | raise RuntimeError("Router socket not created.") |
| 446 | try: |
| 447 | result = msgpack.packb(result) |
| 448 | self.socket.send_multipart([self.req_dict[task_id], b"", result]) |
| 449 | |
| 450 | except Exception as e: |
| 451 | llm_logger.error(f"Send result to zmq client failed: {e}") |
| 452 |
no outgoing calls