| 2641 | } |
| 2642 | |
| 2643 | static int on_next_task(ArrowAsyncDeviceStreamHandler* self, ArrowAsyncTask* task, |
| 2644 | const char* metadata) { |
| 2645 | auto* private_data = reinterpret_cast<PrivateData*>(self->private_data); |
| 2646 | |
| 2647 | if (task == nullptr) { |
| 2648 | std::unique_lock<std::mutex> lock(private_data->state_->mutex_); |
| 2649 | private_data->state_->end_of_stream_ = true; |
| 2650 | lock.unlock(); |
| 2651 | private_data->state_->cv_.notify_one(); |
| 2652 | return 0; |
| 2653 | } |
| 2654 | |
| 2655 | std::shared_ptr<KeyValueMetadata> kvmetadata; |
| 2656 | if (metadata != nullptr) { |
| 2657 | auto maybe_decoded = DecodeMetadata(metadata); |
| 2658 | if (!maybe_decoded.ok()) { |
| 2659 | private_data->state_->error_ = std::move(maybe_decoded).status(); |
| 2660 | private_data->state_->cv_.notify_one(); |
| 2661 | return EINVAL; |
| 2662 | } |
| 2663 | |
| 2664 | kvmetadata = std::move(maybe_decoded->metadata); |
| 2665 | } |
| 2666 | |
| 2667 | std::unique_lock<std::mutex> lock(private_data->state_->mutex_); |
| 2668 | private_data->state_->batches_.push({*task, std::move(kvmetadata)}); |
| 2669 | lock.unlock(); |
| 2670 | private_data->state_->cv_.notify_one(); |
| 2671 | return 0; |
| 2672 | } |
| 2673 | |
| 2674 | static void on_error(ArrowAsyncDeviceStreamHandler* self, int code, const char* message, |
| 2675 | const char* metadata) { |
no test coverage detected