(bytesC chan []byte)
| 424 | } |
| 425 | |
| 426 | func (c *client) decodePipe(bytesC chan []byte) (replyMsg, receivedMsg chan P2PMessage) { |
| 427 | replyMsg = make(chan P2PMessage) |
| 428 | receivedMsg = make(chan P2PMessage) |
| 429 | go func() { |
| 430 | defer close(replyMsg) |
| 431 | defer close(receivedMsg) |
| 432 | for { |
| 433 | var msg P2PMessage |
| 434 | select { |
| 435 | case <-c.ctx.Done(): |
| 436 | return |
| 437 | case bytes, ok := <-bytesC: |
| 438 | if ok { |
| 439 | if len(bytes) == 0 { |
| 440 | continue |
| 441 | } |
| 442 | pa, ptr, err := decodeBytes(bytes, c.verifyFn) |
| 443 | if err != nil { |
| 444 | c.reportError(errors.Errorf("client decodePipe: %w", err)) |
| 445 | continue |
| 446 | } |
| 447 | //TODO: Move to other module |
| 448 | if err := bls.Verify(c.suite, c.remotePubKey, pa.GetAnything().Value, pa.GetSignature()); err != nil { |
| 449 | c.reportError(errors.Errorf("client decodePipe: %w", err)) |
| 450 | continue |
| 451 | } |
| 452 | msg = P2PMessage{Msg: ptr, Sender: pa.GetSender(), RequestNonce: pa.GetRequestNonce()} |
| 453 | if pa.GetReplyFlag() { |
| 454 | select { |
| 455 | case <-c.ctx.Done(): |
| 456 | case replyMsg <- msg: |
| 457 | } |
| 458 | } else { |
| 459 | select { |
| 460 | case <-c.ctx.Done(): |
| 461 | case receivedMsg <- msg: |
| 462 | } |
| 463 | } |
| 464 | } |
| 465 | } |
| 466 | } |
| 467 | }() |
| 468 | return |
| 469 | } |
| 470 | |
| 471 | func (c *client) readPipe() (out chan []byte) { |
| 472 | out = make(chan []byte, 10) |
no test coverage detected