| 394 | std::string_view name; |
| 395 | }; |
| 396 | struct SubmitTask : public Task { |
| 397 | explicit SubmitTask(std::unique_ptr<State> state_holder) |
| 398 | : state_holder(std::move(state_holder)) {} |
| 399 | |
| 400 | struct SubmitTaskCallback { |
| 401 | SubmitTaskCallback(std::unique_ptr<State> state_holder, Future<> task_completion) |
| 402 | : state_holder(std::move(state_holder)), |
| 403 | task_completion(std::move(task_completion)) {} |
| 404 | void operator()(const Result<T>& maybe_item) { |
| 405 | if (!maybe_item.ok()) { |
| 406 | task_completion.MarkFinished(maybe_item.status()); |
| 407 | return; |
| 408 | } |
| 409 | const auto& item = *maybe_item; |
| 410 | if (IsIterationEnd(item)) { |
| 411 | task_completion.MarkFinished(); |
| 412 | return; |
| 413 | } |
| 414 | Status visit_st = state_holder->visitor(item); |
| 415 | if (!visit_st.ok()) { |
| 416 | task_completion.MarkFinished(std::move(visit_st)); |
| 417 | return; |
| 418 | } |
| 419 | state_holder->task_group->AddTask( |
| 420 | std::make_unique<SubmitTask>(std::move(state_holder))); |
| 421 | task_completion.MarkFinished(); |
| 422 | } |
| 423 | std::unique_ptr<State> state_holder; |
| 424 | Future<> task_completion; |
| 425 | }; |
| 426 | |
| 427 | Result<Future<>> operator()() { |
| 428 | Future<> task = Future<>::Make(); |
| 429 | // Consume as many items as we can (those that are already finished) |
| 430 | // synchronously to avoid recursion / stack overflow. |
| 431 | while (true) { |
| 432 | Future<T> next = state_holder->generator(); |
| 433 | if (next.TryAddCallback( |
| 434 | [&] { return SubmitTaskCallback(std::move(state_holder), task); })) { |
| 435 | return task; |
| 436 | } |
| 437 | ARROW_ASSIGN_OR_RAISE(T item, next.result()); |
| 438 | if (IsIterationEnd(item)) { |
| 439 | task.MarkFinished(); |
| 440 | return task; |
| 441 | } |
| 442 | ARROW_RETURN_NOT_OK(state_holder->visitor(item)); |
| 443 | } |
| 444 | } |
| 445 | |
| 446 | std::string_view name() const { return state_holder->name; } |
| 447 | |
| 448 | std::unique_ptr<State> state_holder; |
| 449 | }; |
| 450 | std::unique_ptr<AsyncTaskGroup> task_group = |
| 451 | AsyncTaskGroup::Make(this, [] { return Status::OK(); }); |
| 452 | AsyncTaskGroup* task_group_view = task_group.get(); |
no outgoing calls
no test coverage detected