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

Method ConsumeBuffer

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

Source from the content-addressed store, hash-verified

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

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