| 116 | return result |
| 117 | |
| 118 | def rpc(self, request): |
| 119 | # Match the owner's versioned JSON frame and one connection per RPC. |
| 120 | frame = encoded({"version": 1, "request_id": "perf", "request": request}) |
| 121 | with socket.socket(socket.AF_UNIX, socket.SOCK_STREAM) as stream: |
| 122 | stream.settimeout(45) |
| 123 | stream.connect(self.endpoint) |
| 124 | stream.sendall(struct.pack(">I", len(frame)) + frame) |
| 125 | |
| 126 | def read_exact(size): |
| 127 | data = bytearray() |
| 128 | while len(data) < size: |
| 129 | chunk = stream.recv(size - len(data)) |
| 130 | if not chunk: |
| 131 | raise RuntimeError("owner closed connection before response") |
| 132 | data.extend(chunk) |
| 133 | return data |
| 134 | |
| 135 | size = struct.unpack(">I", read_exact(4))[0] |
| 136 | if not 0 < size <= 8 * 1024 * 1024: |
| 137 | raise RuntimeError(f"invalid owner frame size: {size}") |
| 138 | response = json.loads(read_exact(size))["response"] |
| 139 | if isinstance(response, dict) and "Error" in response: |
| 140 | raise RuntimeError(str(response["Error"])) |
| 141 | return response |
| 142 | |
| 143 | def status(self, turn): |
| 144 | return self.rpc({"TurnStatus": {"session_id": SESSION, "turn_number": turn}})["TurnLifecycle"]["turn"] |