(self, socket_path, debug=False)
| 377 | CMDS = {} |
| 378 | |
| 379 | def __init__(self, socket_path, debug=False): |
| 380 | self.name = f"EngineServer {self.__class__.__name__}" |
| 381 | super().__init__(debug=debug) |
| 382 | self.engine_debug(f"{self.name}: finished setup 1 (_debug={self._engine_debug})") |
| 383 | self.socket_path = socket_path |
| 384 | self.client_id_var = contextvars.ContextVar("client_id", default=None) |
| 385 | # task <--> client id mapping |
| 386 | self.tasks = {} |
| 387 | # child tasks spawned by main tasks |
| 388 | self.child_tasks = {} |
| 389 | self.engine_debug(f"{self.name}: finished setup 2 (_debug={self._engine_debug})") |
| 390 | if self.socket_path is not None: |
| 391 | # create ZeroMQ context |
| 392 | self.context = zmq.asyncio.Context() |
| 393 | # ROUTER socket can handle multiple concurrent requests |
| 394 | self.socket = self.context.socket(zmq.ROUTER) |
| 395 | self.socket.setsockopt(zmq.LINGER, 0) # Discard pending messages immediately disconnect() or close() |
| 396 | self.socket.setsockopt(zmq.SNDHWM, 0) # Unlimited send buffer |
| 397 | self.socket.setsockopt(zmq.RCVHWM, 0) # Unlimited receive buffer |
| 398 | # create socket file |
| 399 | self.socket.bind(f"ipc://{self.socket_path}") |
| 400 | self.engine_debug(f"{self.name}: finished setup 3 (_debug={self._engine_debug})") |
| 401 | |
| 402 | @contextlib.contextmanager |
| 403 | def client_id_context(self, value): |
nothing calls this directly
no test coverage detected