MCPcopy Create free account
hub / github.com/apache/arrow / ConsumeBuffer

Method ConsumeBuffer

cpp/src/arrow/ipc/message.cc:719–763  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

717 }
718
719 Status ConsumeBuffer(std::shared_ptr<Buffer> buffer) {
720 if (buffered_size_ == 0) {
721 while (buffer->size() >= next_required_size_) {
722 auto used_size = next_required_size_;
723 switch (state_) {
724 case State::INITIAL:
725 RETURN_NOT_OK(ConsumeInitialBuffer(buffer));
726 break;
727 case State::METADATA_LENGTH:
728 RETURN_NOT_OK(ConsumeMetadataLengthBuffer(buffer));
729 break;
730 case State::METADATA:
731 if (buffer->size() == next_required_size_) {
732 return ConsumeMetadataBuffer(buffer);
733 } else {
734 auto sliced_buffer = SliceBuffer(buffer, 0, next_required_size_);
735 RETURN_NOT_OK(ConsumeMetadataBuffer(sliced_buffer));
736 }
737 break;
738 case State::BODY:
739 if (buffer->size() == next_required_size_) {
740 return ConsumeBodyBuffer(buffer);
741 } else {
742 auto sliced_buffer = SliceBuffer(buffer, 0, next_required_size_);
743 RETURN_NOT_OK(ConsumeBodyBuffer(sliced_buffer));
744 }
745 break;
746 case State::EOS:
747 return Status::OK();
748 }
749 if (buffer->size() == used_size) {
750 return Status::OK();
751 }
752 buffer = SliceBuffer(buffer, used_size);
753 }
754 }
755
756 if (buffer->size() == 0) {
757 return Status::OK();
758 }
759
760 buffered_size_ += buffer->size();
761 chunks_.push_back(std::move(buffer));
762 return ConsumeChunks();
763 }
764
765 int64_t next_required_size() const { return next_required_size_ - buffered_size_; }
766

Callers 1

ConsumeMethod · 0.45

Calls 4

SliceBufferFunction · 0.85
push_backMethod · 0.80
OKFunction · 0.50
sizeMethod · 0.45

Tested by

no test coverage detected