| 751 | } |
| 752 | |
| 753 | void do_write_request() { |
| 754 | if (dead) |
| 755 | return ; |
| 756 | ASSERT_TRUE(socket.is_connected() > 0); |
| 757 | |
| 758 | while (left > 0 && factory->queue_depth > writings.size() + acking.size()) { |
| 759 | StressFactory::Message *m = new StressFactory::Message( |
| 760 | factory->rs, ++index, |
| 761 | factory->rd() % factory->max_message_length); |
| 762 | std::cerr << " client " << this << " generate message " << m->idx << " length " << m->len << std::endl; |
| 763 | ASSERT_EQ(m->len, m->content.size()); |
| 764 | writings.push_back(m); |
| 765 | --left; |
| 766 | --factory->message_left; |
| 767 | } |
| 768 | |
| 769 | while (!writings.empty()) { |
| 770 | StressFactory::Message *m = writings.front(); |
| 771 | bufferlist bl; |
| 772 | bl.append(m->content.data() + write_offset, m->content.size() - write_offset); |
| 773 | ssize_t r = socket.send(bl, false); |
| 774 | if (r == 0) |
| 775 | break; |
| 776 | std::cerr << " client " << this << " send " << m->idx << " len " << r << " content: " << std::endl; |
| 777 | ASSERT_TRUE(r >= 0); |
| 778 | write_offset += r; |
| 779 | if (write_offset == m->content.size()) { |
| 780 | write_offset = 0; |
| 781 | writings.pop_front(); |
| 782 | acking.push_back(m); |
| 783 | } |
| 784 | } |
| 785 | if (writings.empty() && write_enabled) { |
| 786 | center->delete_file_event(socket.fd(), EVENT_WRITABLE); |
| 787 | write_enabled = false; |
| 788 | } else if (!writings.empty() && !write_enabled) { |
| 789 | ASSERT_EQ(0, center->create_file_event( |
| 790 | socket.fd(), EVENT_WRITABLE, &write_ctxt)); |
| 791 | write_enabled = true; |
| 792 | } |
| 793 | } |
| 794 | |
| 795 | bool finish() const { |
| 796 | return left == 0 && acking.empty() && writings.empty(); |
no test coverage detected