| 137 | stream_finished_(false) {} |
| 138 | |
| 139 | ::arrow::Result<std::unique_ptr<ipc::Message>> ReadNextMessage() override { |
| 140 | if (stream_finished_) { |
| 141 | return nullptr; |
| 142 | } |
| 143 | internal::FlightData* data; |
| 144 | peekable_reader_->Next(&data); |
| 145 | if (!data) { |
| 146 | stream_finished_ = true; |
| 147 | ARROW_RETURN_NOT_OK(stream_->Finish(Status::OK())); |
| 148 | return nullptr; |
| 149 | } |
| 150 | if (data->body) { |
| 151 | ARROW_ASSIGN_OR_RAISE(data->body, Buffer::ViewOrCopy(data->body, memory_manager_)); |
| 152 | } |
| 153 | // Validate IPC message |
| 154 | auto result = data->OpenMessage(); |
| 155 | if (!result.ok()) { |
| 156 | stream_finished_ = true; |
| 157 | ARROW_RETURN_NOT_OK(stream_->Finish(std::move(result).status())); |
| 158 | return nullptr; |
| 159 | } |
| 160 | *app_metadata_ = std::move(data->app_metadata); |
| 161 | return result; |
| 162 | } |
| 163 | |
| 164 | private: |
| 165 | std::shared_ptr<internal::ClientDataStream> stream_; |
no test coverage detected