MCPcopy Create free account
hub / github.com/ceph/ceph / do_read_request

Method do_read_request

src/test/msgr/test_async_networkstack.cc:704–751  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

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)

Callers 2

do_requestMethod · 0.45
do_requestMethod · 0.45

Calls 10

is_connectedMethod · 0.45
emptyMethod · 0.45
frontMethod · 0.45
resizeMethod · 0.45
readMethod · 0.45
dataMethod · 0.45
verifyMethod · 0.45
pop_frontMethod · 0.45
clearMethod · 0.45

Tested by

no test coverage detected