Callback when a message is received. Parameters ---------- message : bytearray The bytes received
(self, message)
| 211 | self._init_req_nbytes = 0 |
| 212 | |
| 213 | def on_message(self, message): |
| 214 | """Callback when a message is received. |
| 215 | |
| 216 | Parameters |
| 217 | ---------- |
| 218 | message : bytearray |
| 219 | The bytes received |
| 220 | """ |
| 221 | assert isinstance(message, bytes) |
| 222 | if self._init_req_nbytes: |
| 223 | self._init_conn(message) |
| 224 | return |
| 225 | |
| 226 | self._data += message |
| 227 | |
| 228 | while True: |
| 229 | if self._msg_size == 0: |
| 230 | if len(self._data) >= 4: |
| 231 | self._msg_size = struct.unpack("<i", self._data[:4])[0] |
| 232 | if self._msg_size <= 0 or self._msg_size > MAX_TRACKER_MSG_BYTES: |
| 233 | logger.warning( |
| 234 | "Invalid msg_size %d from %s; closing connection", |
| 235 | self._msg_size, |
| 236 | self.name(), |
| 237 | ) |
| 238 | self.close() |
| 239 | return |
| 240 | del self._data[:4] |
| 241 | else: |
| 242 | return |
| 243 | if self._msg_size != 0 and len(self._data) >= self._msg_size: |
| 244 | msg = (bytes(self._data[: self._msg_size])).decode("utf-8") |
| 245 | del self._data[: self._msg_size] |
| 246 | self._msg_size = 0 |
| 247 | try: |
| 248 | self.call_handler(json.loads(msg)) |
| 249 | except Exception: # pylint: disable=broad-except |
| 250 | logger.warning("Error handling message from %s", self.name(), exc_info=True) |
| 251 | self.close() |
| 252 | return |
| 253 | else: |
| 254 | return |
| 255 | |
| 256 | def ret_value(self, data): |
| 257 | """return value to the output""" |
nothing calls this directly
no test coverage detected