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