| 702 | } |
| 703 | |
| 704 | void do_read_request() { |
| 705 | if (dead) |
| 706 | return ; |
| 707 | ASSERT_TRUE(socket.is_connected() >= 0); |
| 708 | if (!socket.is_connected()) |
| 709 | return ; |
| 710 | ASSERT_TRUE(!acking.empty() || first); |
| 711 | if (first) { |
| 712 | first = false; |
| 713 | center->dispatch_event_external(&write_ctxt); |
| 714 | if (acking.empty()) |
| 715 | return ; |
| 716 | } |
| 717 | StressFactory::Message *m = acking.front(); |
| 718 | int r = 0; |
| 719 | if (buffer.empty()) |
| 720 | buffer.resize(m->len); |
| 721 | bool must_no = false; |
| 722 | while (true) { |
| 723 | r = socket.read((char*)buffer.data() + read_offset, |
| 724 | m->len - read_offset); |
| 725 | ASSERT_TRUE(r == -EAGAIN || r > 0); |
| 726 | if (r == -EAGAIN) |
| 727 | break; |
| 728 | read_offset += r; |
| 729 | |
| 730 | std::cerr << " client " << this << " receive " << m->idx << " len " << r << " content: " << std::endl; |
| 731 | ASSERT_FALSE(must_no); |
| 732 | if ((m->len - read_offset) == 0) { |
| 733 | ASSERT_TRUE(m->verify(buffer.data(), 0)); |
| 734 | delete m; |
| 735 | acking.pop_front(); |
| 736 | read_offset = 0; |
| 737 | buffer.clear(); |
| 738 | if (acking.empty()) { |
| 739 | m = &homeless_message; |
| 740 | must_no = true; |
| 741 | } else { |
| 742 | m = acking.front(); |
| 743 | buffer.resize(m->len); |
| 744 | } |
| 745 | } |
| 746 | } |
| 747 | if (acking.empty()) { |
| 748 | center->dispatch_event_external(&write_ctxt); |
| 749 | return ; |
| 750 | } |
| 751 | } |
| 752 | |
| 753 | void do_write_request() { |
| 754 | if (dead) |
no test coverage detected