| 295 | await self._pub_socket.send(msg.serialize()) |
| 296 | |
| 297 | async def sub(self, callback: Callable[[Message], Any]): |
| 298 | # Subscribe and handle messages |
| 299 | if not self._sub_socket: |
| 300 | ep = f"{self.endpoint.protocol}://{self.endpoint.address}" |
| 301 | self._sub_socket = self.context.socket(zmq.SUB) |
| 302 | self._sub_socket.connect(ep) |
| 303 | self._sub_socket.setsockopt_string(zmq.SUBSCRIBE, "") |
| 304 | |
| 305 | async def loop(): |
| 306 | while True: |
| 307 | raw = await self._sub_socket.recv() |
| 308 | msg = Message.deserialize(raw) |
| 309 | result = callback(msg) |
| 310 | if asyncio.iscoroutine(result): |
| 311 | await result |
| 312 | |
| 313 | self._task = asyncio.create_task(loop()) |
| 314 | |
| 315 | |
| 316 | # ========================== |