| 110 | } |
| 111 | |
| 112 | arrow::Result<FlightStreamChunk> Next() override { |
| 113 | FlightStreamChunk out; |
| 114 | internal::FlightData* data = nullptr; |
| 115 | peekable_reader_->Peek(&data); |
| 116 | if (!data) { |
| 117 | out.app_metadata = nullptr; |
| 118 | out.data = nullptr; |
| 119 | return out; |
| 120 | } |
| 121 | |
| 122 | if (!data->metadata) { |
| 123 | // Metadata-only (data->metadata is the IPC header) |
| 124 | out.app_metadata = data->app_metadata; |
| 125 | out.data = nullptr; |
| 126 | peekable_reader_->Next(&data); |
| 127 | return out; |
| 128 | } |
| 129 | |
| 130 | if (!batch_reader_) { |
| 131 | RETURN_NOT_OK(EnsureDataStarted()); |
| 132 | // re-peek here since EnsureDataStarted() advances the stream |
| 133 | return Next(); |
| 134 | } |
| 135 | RETURN_NOT_OK(batch_reader_->ReadNext(&out.data)); |
| 136 | out.app_metadata = std::move(app_metadata_); |
| 137 | return out; |
| 138 | } |
| 139 | |
| 140 | ipc::ReadStats stats() const override { |
| 141 | if (batch_reader_ == nullptr) { |
no test coverage detected