Try to read P2P messages from the recv buffer. This method reads data from the buffer in a loop. It deserializes, parses and verifies the P2P header, then passes the P2P payload to the on_message callback for processing.
(self)
| 216 | self._on_data() |
| 217 | |
| 218 | def _on_data(self): |
| 219 | """Try to read P2P messages from the recv buffer. |
| 220 | |
| 221 | This method reads data from the buffer in a loop. It deserializes, |
| 222 | parses and verifies the P2P header, then passes the P2P payload to |
| 223 | the on_message callback for processing.""" |
| 224 | try: |
| 225 | while True: |
| 226 | if len(self.recvbuf) < 4: |
| 227 | return |
| 228 | if self.recvbuf[:4] != self.magic_bytes: |
| 229 | raise ValueError("magic bytes mismatch: {} != {}".format(repr(self.magic_bytes), repr(self.recvbuf))) |
| 230 | if len(self.recvbuf) < 4 + 12 + 4 + 4: |
| 231 | return |
| 232 | msgtype = self.recvbuf[4:4+12].split(b"\x00", 1)[0] |
| 233 | msglen = struct.unpack("<i", self.recvbuf[4+12:4+12+4])[0] |
| 234 | checksum = self.recvbuf[4+12+4:4+12+4+4] |
| 235 | if len(self.recvbuf) < 4 + 12 + 4 + 4 + msglen: |
| 236 | return |
| 237 | msg = self.recvbuf[4+12+4+4:4+12+4+4+msglen] |
| 238 | th = sha256(msg) |
| 239 | h = sha256(th) |
| 240 | if checksum != h[:4]: |
| 241 | raise ValueError("got bad checksum " + repr(self.recvbuf)) |
| 242 | self.recvbuf = self.recvbuf[4+12+4+4+msglen:] |
| 243 | if msgtype not in MESSAGEMAP: |
| 244 | raise ValueError("Received unknown msgtype from %s:%d: '%s' %s" % (self.dstaddr, self.dstport, msgtype, repr(msg))) |
| 245 | f = BytesIO(msg) |
| 246 | t = MESSAGEMAP[msgtype]() |
| 247 | t.deserialize(f) |
| 248 | self._log_message("receive", t) |
| 249 | self.on_message(t) |
| 250 | except Exception as e: |
| 251 | logger.exception('Error reading message:', repr(e)) |
| 252 | raise |
| 253 | |
| 254 | def on_message(self, message): |
| 255 | """Callback for processing a P2P payload. Must be overridden by derived class.""" |
no test coverage detected