TLA Line data Source code
1 : //
2 : // Copyright (c) 2026 Michael Vandeberg
3 : // Copyright (c) 2026 Steve Gerbino
4 : //
5 : // Distributed under the Boost Software License, Version 1.0. (See accompanying
6 : // file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
7 : //
8 : // Official repository: https://github.com/cppalliance/capy
9 : //
10 :
11 : #ifndef BOOST_CAPY_WHEN_ANY_HPP
12 : #define BOOST_CAPY_WHEN_ANY_HPP
13 :
14 : #include <boost/capy/detail/config.hpp>
15 : #include <boost/capy/detail/io_result_combinators.hpp>
16 : #include <boost/capy/continuation.hpp>
17 : #include <boost/capy/concept/executor.hpp>
18 : #include <boost/capy/concept/io_awaitable.hpp>
19 : #include <coroutine>
20 : #include <boost/capy/ex/executor_ref.hpp>
21 : #include <boost/capy/ex/frame_alloc_mixin.hpp>
22 : #include <boost/capy/ex/frame_allocator.hpp>
23 : #include <boost/capy/ex/io_env.hpp>
24 : #include <boost/capy/task.hpp>
25 :
26 : #include <array>
27 : #include <atomic>
28 : #include <exception>
29 : #include <memory>
30 : #include <mutex>
31 : #include <optional>
32 : #include <ranges>
33 : #include <stdexcept>
34 : #include <stop_token>
35 : #include <tuple>
36 : #include <type_traits>
37 : #include <utility>
38 : #include <variant>
39 : #include <vector>
40 :
41 : /*
42 : when_any - Race multiple io_result tasks, select first success
43 : =============================================================
44 :
45 : OVERVIEW:
46 : ---------
47 : when_any launches N io_result-returning tasks concurrently. A task
48 : wins by returning !ec; errors and exceptions do not win. Once a
49 : winner is found, stop is requested for siblings and the winner's
50 : payload is returned. If no winner exists (all fail), one of the
51 : failures is surfaced (an error_code at variant index 0, or a child's
52 : exception rethrown); which one is unspecified.
53 :
54 : ARCHITECTURE:
55 : -------------
56 : The design mirrors when_all but with inverted completion semantics:
57 :
58 : when_all: complete when remaining_count reaches 0 (all done)
59 : when_any: complete when has_winner becomes true (first done)
60 : BUT still wait for remaining_count to reach 0 for cleanup
61 :
62 : Key components:
63 : - when_any_core: Shared state tracking winner and completion
64 : - when_any_io_runner: Wrapper coroutine for each child task
65 : - when_any_io_launcher/when_any_io_homogeneous_launcher:
66 : Awaitables that start all runners concurrently
67 :
68 : CRITICAL INVARIANTS:
69 : --------------------
70 : 1. Only a task returning !ec can become the winner (via atomic CAS)
71 : 2. All tasks must complete before parent resumes (cleanup safety)
72 : 3. Stop is requested immediately when winner is determined
73 : 4. Exceptions and errors do not claim winner status
74 :
75 : POSITIONAL VARIANT:
76 : -------------------
77 : The variadic overload returns std::variant<error_code, R1, R2, ..., Rn>.
78 : Index 0 is error_code (failure/no-winner). Index 1..N identifies the
79 : winning child and carries its payload.
80 :
81 : RANGE OVERLOAD:
82 : ---------------
83 : The range overload returns variant<error_code, pair<size_t, T>> for
84 : non-void children or variant<error_code, size_t> for void children.
85 :
86 : MEMORY MODEL:
87 : -------------
88 : Synchronization chain from winner's write to parent's read:
89 :
90 : 1. Winner thread writes result_ (non-atomic)
91 : 2. Winner thread calls signal_completion() -> fetch_sub(acq_rel) on remaining_count_
92 : 3. Last task thread (may be winner or non-winner) calls signal_completion()
93 : -> fetch_sub(acq_rel) on remaining_count_, observing count becomes 0
94 : 4. Last task returns caller_ex_.dispatch(continuation_) via symmetric transfer
95 : 5. Parent coroutine resumes and reads result_
96 :
97 : Synchronization analysis:
98 : - All fetch_sub operations on remaining_count_ form a release sequence
99 : - Winner's fetch_sub releases; subsequent fetch_sub operations participate
100 : in the modification order of remaining_count_
101 : - Last task's fetch_sub(acq_rel) synchronizes-with prior releases in the
102 : modification order, establishing happens-before from winner's writes
103 : - Executor dispatch() is expected to provide queue-based synchronization
104 : (release-on-post, acquire-on-execute) completing the chain to parent
105 : - Even inline executors work (same thread = sequenced-before)
106 :
107 : EXCEPTION SEMANTICS:
108 : --------------------
109 : Exceptions do NOT claim winner status. If a child throws, the exception
110 : is recorded but the combinator keeps waiting for a success. Only when
111 : all children complete without a winner is a failure surfaced. There is
112 : no priority between errors and exceptions, and no guarantee about which
113 : child's failure is reported: the result either returns an error_code at
114 : variant index 0 or rethrows a child's exception.
115 : */
116 :
117 : namespace boost {
118 : namespace capy {
119 :
120 : namespace detail {
121 :
122 : /** Core shared state for when_any operations.
123 :
124 : Contains all members and methods common to both heterogeneous (variadic)
125 : and homogeneous (range) when_any implementations. State classes embed
126 : this via composition to avoid CRTP destructor ordering issues.
127 :
128 : @par Thread Safety
129 : Atomic operations protect winner selection and completion count.
130 : */
131 : struct when_any_core
132 : {
133 : std::atomic<std::size_t> remaining_count_;
134 : std::size_t winner_index_{0};
135 : std::exception_ptr winner_exception_;
136 : std::stop_source stop_source_;
137 :
138 : // Bridges parent's stop token to our stop_source
139 : struct stop_callback_fn
140 : {
141 : std::stop_source* source_;
142 HIT 3 : void operator()() const noexcept { source_->request_stop(); }
143 : };
144 : using stop_callback_t = std::stop_callback<stop_callback_fn>;
145 : std::optional<stop_callback_t> parent_stop_callback_;
146 :
147 : continuation continuation_;
148 : io_env const* caller_env_ = nullptr;
149 :
150 : // Placed last to avoid padding (1-byte atomic followed by 8-byte aligned members)
151 : std::atomic<bool> has_winner_{false};
152 :
153 39 : explicit when_any_core(std::size_t count) noexcept
154 39 : : remaining_count_(count)
155 : {
156 39 : }
157 :
158 : /** Atomically claim winner status; exactly one task succeeds. */
159 60 : bool try_win(std::size_t index) noexcept
160 : {
161 60 : bool expected = false;
162 60 : if(has_winner_.compare_exchange_strong(
163 : expected, true, std::memory_order_acq_rel))
164 : {
165 28 : winner_index_ = index;
166 28 : stop_source_.request_stop();
167 28 : return true;
168 : }
169 32 : return false;
170 : }
171 :
172 : /** @pre try_win() returned true. */
173 1 : void set_winner_exception(std::exception_ptr ep) noexcept
174 : {
175 1 : winner_exception_ = ep;
176 1 : }
177 :
178 : // Runners signal completion directly via final_suspend; no member function needed.
179 : };
180 :
181 : } // namespace detail
182 :
183 : namespace detail {
184 :
185 : // State for io_result-aware when_any: only !ec wins.
186 : template<typename... Ts>
187 : struct when_any_io_state
188 : {
189 : static constexpr std::size_t task_count = sizeof...(Ts);
190 : using variant_type = std::variant<std::error_code, Ts...>;
191 :
192 : when_any_core core_;
193 : std::optional<variant_type> result_;
194 : std::array<continuation, task_count> runner_handles_{};
195 :
196 : // A failure (error or exception) for the all-fail case. record_error
197 : // and record_exception overwrite each other, so which one survives is
198 : // unspecified (no priority between errors and exceptions).
199 : std::mutex failure_mu_;
200 : std::error_code last_error_;
201 : std::exception_ptr last_exception_;
202 :
203 23 : when_any_io_state()
204 23 : : core_(task_count)
205 : {
206 23 : }
207 :
208 15 : void record_error(std::error_code ec)
209 : {
210 15 : std::lock_guard lk(failure_mu_);
211 15 : last_error_ = ec;
212 15 : last_exception_ = nullptr;
213 15 : }
214 :
215 7 : void record_exception(std::exception_ptr ep)
216 : {
217 7 : std::lock_guard lk(failure_mu_);
218 7 : last_exception_ = ep;
219 7 : last_error_ = {};
220 7 : }
221 : };
222 :
223 : // Wrapper coroutine for io_result-aware when_any children.
224 : // unhandled_exception records the exception but does NOT claim winner status.
225 : template<typename StateType>
226 : struct BOOST_CAPY_CORO_DESTROY_WHEN_COMPLETE when_any_io_runner
227 : {
228 : struct promise_type
229 : : frame_alloc_mixin
230 : {
231 : StateType* state_ = nullptr;
232 : std::size_t index_ = 0;
233 : io_env env_;
234 :
235 95 : when_any_io_runner get_return_object() noexcept
236 : {
237 : return when_any_io_runner(
238 95 : std::coroutine_handle<promise_type>::from_promise(*this));
239 : }
240 :
241 95 : std::suspend_always initial_suspend() noexcept { return {}; }
242 :
243 95 : auto final_suspend() noexcept
244 : {
245 : struct awaiter
246 : {
247 : promise_type* p_;
248 95 : bool await_ready() const noexcept { return false; }
249 95 : auto await_suspend(std::coroutine_handle<> h) noexcept
250 : {
251 95 : auto& core = p_->state_->core_;
252 95 : auto* counter = &core.remaining_count_;
253 95 : auto* caller_env = core.caller_env_;
254 95 : auto& cont = core.continuation_;
255 :
256 95 : h.destroy();
257 :
258 95 : auto remaining = counter->fetch_sub(1, std::memory_order_acq_rel);
259 95 : if(remaining == 1)
260 39 : return detail::symmetric_transfer(caller_env->executor.dispatch(cont));
261 56 : return detail::symmetric_transfer(std::noop_coroutine());
262 : }
263 : void await_resume() const noexcept {} // LCOV_EXCL_LINE final_suspend awaiter, never resumed
264 : };
265 95 : return awaiter{this};
266 : }
267 :
268 82 : void return_void() noexcept {}
269 :
270 : // Exceptions do NOT win in io_result when_any
271 13 : void unhandled_exception() noexcept
272 : {
273 13 : state_->record_exception(std::current_exception());
274 13 : }
275 :
276 : template<class Awaitable>
277 : struct transform_awaiter
278 : {
279 : std::decay_t<Awaitable> a_;
280 : promise_type* p_;
281 :
282 95 : bool await_ready() { return a_.await_ready(); }
283 95 : decltype(auto) await_resume() { return a_.await_resume(); }
284 :
285 : template<class Promise>
286 94 : auto await_suspend(std::coroutine_handle<Promise> h)
287 : {
288 : using R = decltype(a_.await_suspend(h, &p_->env_));
289 : if constexpr (std::is_same_v<R, std::coroutine_handle<>>)
290 94 : return detail::symmetric_transfer(a_.await_suspend(h, &p_->env_));
291 : else
292 : return a_.await_suspend(h, &p_->env_);
293 : }
294 : };
295 :
296 : template<class Awaitable>
297 95 : auto await_transform(Awaitable&& a)
298 : {
299 : using A = std::decay_t<Awaitable>;
300 : if constexpr (IoAwaitable<A>)
301 : {
302 : return transform_awaiter<Awaitable>{
303 188 : std::forward<Awaitable>(a), this};
304 : }
305 : else
306 : {
307 : static_assert(sizeof(A) == 0, "requires IoAwaitable");
308 : }
309 93 : }
310 : };
311 :
312 : std::coroutine_handle<promise_type> h_;
313 :
314 95 : explicit when_any_io_runner(std::coroutine_handle<promise_type> h) noexcept
315 95 : : h_(h)
316 : {
317 95 : }
318 :
319 : when_any_io_runner(when_any_io_runner&& other) noexcept
320 : : h_(std::exchange(other.h_, nullptr))
321 : {
322 : }
323 :
324 : when_any_io_runner(when_any_io_runner const&) = delete;
325 : when_any_io_runner& operator=(when_any_io_runner const&) = delete;
326 : when_any_io_runner& operator=(when_any_io_runner&&) = delete;
327 :
328 95 : auto release() noexcept
329 : {
330 95 : return std::exchange(h_, nullptr);
331 : }
332 : };
333 :
334 : // Runner coroutine: only tries to win when the child returns !ec.
335 : template<std::size_t I, IoAwaitable Awaitable, typename StateType>
336 : when_any_io_runner<StateType>
337 41 : make_when_any_io_runner(Awaitable inner, StateType* state)
338 : {
339 : auto result = co_await std::move(inner);
340 :
341 : if(!result.ec)
342 : {
343 : // Success: try to claim winner
344 : if(state->core_.try_win(I))
345 : {
346 : try
347 : {
348 : state->result_.emplace(
349 : std::in_place_index<I + 1>,
350 : detail::extract_io_payload(std::move(result)));
351 : }
352 : catch(...)
353 : {
354 : state->core_.set_winner_exception(std::current_exception());
355 : }
356 : }
357 : }
358 : else
359 : {
360 : // Error: record but don't win
361 : state->record_error(result.ec);
362 : }
363 82 : }
364 :
365 : // Launcher for io_result-aware when_any.
366 : template<IoAwaitable... Awaitables>
367 : class when_any_io_launcher
368 : {
369 : using state_type = when_any_io_state<
370 : io_result_payload_t<awaitable_result_t<Awaitables>>...>;
371 :
372 : std::tuple<Awaitables...>* tasks_;
373 : state_type* state_;
374 :
375 : public:
376 23 : when_any_io_launcher(
377 : std::tuple<Awaitables...>* tasks,
378 : state_type* state)
379 23 : : tasks_(tasks)
380 23 : , state_(state)
381 : {
382 23 : }
383 :
384 23 : bool await_ready() const noexcept
385 : {
386 23 : return sizeof...(Awaitables) == 0;
387 : }
388 :
389 23 : std::coroutine_handle<> await_suspend(
390 : std::coroutine_handle<> continuation, io_env const* caller_env)
391 : {
392 23 : state_->core_.continuation_.h = continuation;
393 23 : state_->core_.caller_env_ = caller_env;
394 :
395 23 : if(caller_env->stop_token.stop_possible())
396 : {
397 4 : state_->core_.parent_stop_callback_.emplace(
398 2 : caller_env->stop_token,
399 2 : when_any_core::stop_callback_fn{&state_->core_.stop_source_});
400 :
401 2 : if(caller_env->stop_token.stop_requested())
402 1 : state_->core_.stop_source_.request_stop();
403 : }
404 :
405 23 : auto token = state_->core_.stop_source_.get_token();
406 23 : launch_all(std::index_sequence_for<Awaitables...>{},
407 : caller_env->executor, token);
408 :
409 46 : return std::noop_coroutine();
410 23 : }
411 :
412 23 : void await_resume() const noexcept {}
413 :
414 : private:
415 : template<std::size_t... Is>
416 23 : void launch_all(std::index_sequence<Is...>,
417 : executor_ref ex, std::stop_token token)
418 : {
419 23 : (..., launch_one<Is>(ex, token));
420 23 : }
421 :
422 : template<std::size_t I>
423 41 : void launch_one(executor_ref caller_ex, std::stop_token token)
424 : {
425 41 : auto runner = make_when_any_io_runner<I>(
426 41 : std::move(std::get<I>(*tasks_)), state_);
427 :
428 41 : auto h = runner.release();
429 41 : h.promise().state_ = state_;
430 41 : h.promise().index_ = I;
431 41 : h.promise().env_ = io_env{caller_ex, token,
432 41 : state_->core_.caller_env_->frame_allocator};
433 :
434 41 : state_->runner_handles_[I].h = std::coroutine_handle<>{h};
435 41 : caller_ex.post(state_->runner_handles_[I]);
436 82 : }
437 : };
438 :
439 : /** Shared state for homogeneous io_result-aware when_any (range overload).
440 :
441 : @tparam T The payload type extracted from io_result.
442 : */
443 : template<typename T>
444 : struct when_any_io_homogeneous_state
445 : {
446 : when_any_core core_;
447 : std::optional<T> result_;
448 : std::unique_ptr<continuation[]> runner_handles_;
449 :
450 : std::mutex failure_mu_;
451 : std::error_code last_error_;
452 : std::exception_ptr last_exception_;
453 :
454 13 : explicit when_any_io_homogeneous_state(std::size_t count)
455 13 : : core_(count)
456 13 : , runner_handles_(std::make_unique<continuation[]>(count))
457 : {
458 13 : }
459 :
460 6 : void record_error(std::error_code ec)
461 : {
462 6 : std::lock_guard lk(failure_mu_);
463 6 : last_error_ = ec;
464 6 : last_exception_ = nullptr;
465 6 : }
466 :
467 4 : void record_exception(std::exception_ptr ep)
468 : {
469 4 : std::lock_guard lk(failure_mu_);
470 4 : last_exception_ = ep;
471 4 : last_error_ = {};
472 4 : }
473 : };
474 :
475 : /** Specialization for void io_result children (no payload storage). */
476 : template<>
477 : struct when_any_io_homogeneous_state<std::tuple<>>
478 : {
479 : when_any_core core_;
480 : std::unique_ptr<continuation[]> runner_handles_;
481 :
482 : std::mutex failure_mu_;
483 : std::error_code last_error_;
484 : std::exception_ptr last_exception_;
485 :
486 3 : explicit when_any_io_homogeneous_state(std::size_t count)
487 3 : : core_(count)
488 3 : , runner_handles_(std::make_unique<continuation[]>(count))
489 : {
490 3 : }
491 :
492 1 : void record_error(std::error_code ec)
493 : {
494 1 : std::lock_guard lk(failure_mu_);
495 1 : last_error_ = ec;
496 1 : last_exception_ = nullptr;
497 1 : }
498 :
499 2 : void record_exception(std::exception_ptr ep)
500 : {
501 2 : std::lock_guard lk(failure_mu_);
502 2 : last_exception_ = ep;
503 2 : last_error_ = {};
504 2 : }
505 : };
506 :
507 : /** Create an io_result-aware runner for homogeneous when_any (range path).
508 :
509 : Only tries to win when the child returns !ec.
510 : */
511 : template<IoAwaitable Awaitable, typename StateType>
512 : when_any_io_runner<StateType>
513 54 : make_when_any_io_homogeneous_runner(
514 : Awaitable inner, StateType* state, std::size_t index)
515 : {
516 : auto result = co_await std::move(inner);
517 :
518 : if(!result.ec)
519 : {
520 : if(state->core_.try_win(index))
521 : {
522 : using PayloadT = io_result_payload_t<
523 : awaitable_result_t<Awaitable>>;
524 : if constexpr (!std::is_same_v<PayloadT, std::tuple<>>)
525 : {
526 : try
527 : {
528 : state->result_.emplace(
529 : extract_io_payload(std::move(result)));
530 : }
531 : catch(...)
532 : {
533 : state->core_.set_winner_exception(
534 : std::current_exception());
535 : }
536 : }
537 : }
538 : }
539 : else
540 : {
541 : state->record_error(result.ec);
542 : }
543 108 : }
544 :
545 : /** Starts all io_result-aware homogeneous runners concurrently. */
546 : template<IoAwaitableRange Range>
547 : class when_any_io_homogeneous_launcher
548 : {
549 : using Awaitable = std::ranges::range_value_t<Range>;
550 : using PayloadT = io_result_payload_t<awaitable_result_t<Awaitable>>;
551 :
552 : Range* range_;
553 : when_any_io_homogeneous_state<PayloadT>* state_;
554 :
555 : public:
556 16 : when_any_io_homogeneous_launcher(
557 : Range* range,
558 : when_any_io_homogeneous_state<PayloadT>* state)
559 16 : : range_(range)
560 16 : , state_(state)
561 : {
562 16 : }
563 :
564 16 : bool await_ready() const noexcept
565 : {
566 16 : return std::ranges::empty(*range_);
567 : }
568 :
569 16 : std::coroutine_handle<> await_suspend(
570 : std::coroutine_handle<> continuation, io_env const* caller_env)
571 : {
572 16 : state_->core_.continuation_.h = continuation;
573 16 : state_->core_.caller_env_ = caller_env;
574 :
575 16 : if(caller_env->stop_token.stop_possible())
576 : {
577 4 : state_->core_.parent_stop_callback_.emplace(
578 2 : caller_env->stop_token,
579 2 : when_any_core::stop_callback_fn{&state_->core_.stop_source_});
580 :
581 2 : if(caller_env->stop_token.stop_requested())
582 1 : state_->core_.stop_source_.request_stop();
583 : }
584 :
585 16 : auto token = state_->core_.stop_source_.get_token();
586 :
587 : // Phase 1: Create all runners without dispatching.
588 16 : std::size_t index = 0;
589 70 : for(auto&& a : *range_)
590 : {
591 54 : auto runner = make_when_any_io_homogeneous_runner(
592 54 : std::move(a), state_, index);
593 :
594 54 : auto h = runner.release();
595 54 : h.promise().state_ = state_;
596 54 : h.promise().index_ = index;
597 54 : h.promise().env_ = io_env{caller_env->executor, token,
598 54 : caller_env->frame_allocator};
599 :
600 54 : state_->runner_handles_[index].h = std::coroutine_handle<>{h};
601 54 : ++index;
602 : }
603 :
604 : // Phase 2: Post all runners. Any may complete synchronously.
605 16 : auto* handles = state_->runner_handles_.get();
606 16 : std::size_t count = state_->core_.remaining_count_.load(std::memory_order_relaxed);
607 70 : for(std::size_t i = 0; i < count; ++i)
608 54 : caller_env->executor.post(handles[i]);
609 :
610 32 : return std::noop_coroutine();
611 70 : }
612 :
613 16 : void await_resume() const noexcept {}
614 : };
615 :
616 : } // namespace detail
617 :
618 : /** Race a range of io_result-returning awaitables (non-void payloads).
619 :
620 : Only a child returning !ec can win. Errors and exceptions do not
621 : claim winner status. If all children fail, an unspecified one of
622 : the failures is reported — either an error_code at variant index 0,
623 : or a child's exception rethrown.
624 :
625 : @par Await-effects
626 :
627 : Takes ownership of the range, creates one wrapper coroutine per
628 : element, then posts every wrapper to the caller's executor. All
629 : children therefore run concurrently, each awaited with the caller's
630 : executor and frame allocator and with a stop token owned by this
631 : operation.
632 :
633 : Awaiting an empty range throws `std::invalid_argument` before any
634 : child is started.
635 :
636 : The first child to await-return a zero `ec` claims the win. Claiming
637 : the win requests stop on the operation's own stop token, which every
638 : sibling observes through the stop token it was awaited with. A child
639 : that await-returns a non-zero `ec`, or that exits via an exception,
640 : does not claim the win and does not request stop. The operation keeps
641 : waiting for a success. A stop request on the caller's stop token is
642 : also forwarded to every child.
643 :
644 : The await completes only after every child has finished, regardless of
645 : whether a win was claimed.
646 :
647 : @par Await-returns
648 : An object of type
649 : `std::variant<std::error_code, std::pair<std::size_t, PayloadT>>`,
650 : where `PayloadT` is the payload of one child's `io_result`.
651 :
652 : @li Index 1 holds the winner's position in the input range paired
653 : with its payload.
654 : @li Index 0 holds a non-zero `error_code` when no child won, that is,
655 : when every child failed. It is the `ec` of one of the failed
656 : children; which one is unspecified.
657 :
658 : A child that succeeds after the win has already been claimed
659 : contributes nothing: its payload is discarded.
660 :
661 : If no child won and the failure selected for reporting is an
662 : exception rather than an `ec`, that exception is rethrown instead of
663 : await-returning. The choice of child is unspecified.
664 :
665 : @par Await-postcondition
666 : Every child has finished. If at least one child await-returned a zero
667 : `ec`, the result holds index 1, unless producing the winner's payload
668 : threw, in which case that exception is rethrown. Otherwise the result
669 : holds index 0, or a failed child's exception is rethrown.
670 :
671 : @par Remarks
672 : Supports _IoAwaitable cancellation_. A canceled child await-returns a
673 : non-zero `ec` and so cannot win; if no child has already succeeded,
674 : the result settles at index 0.
675 :
676 : @par Thread Safety
677 : The returned task must be awaited from a single execution context.
678 : Child awaitables execute concurrently but complete through the caller's
679 : executor.
680 :
681 : @param awaitables Range of io_result-returning awaitables (must
682 : not be empty).
683 :
684 : @return A task yielding variant<error_code, pair<size_t, PayloadT>>
685 : where index 0 is failure and index 1 carries the winner's
686 : index and payload.
687 :
688 : @throws std::invalid_argument if range is empty.
689 :
690 : @par Exception Safety
691 : The winner's exception is rethrown if extracting or
692 : move-constructing the winning payload throws. In that case a winner
693 : was found, but its result could not be produced. If all children
694 : fail and the reported failure is an exception, that child's
695 : exception is rethrown (which child is unspecified).
696 :
697 : @par Example
698 : @code
699 : task<void> example()
700 : {
701 : std::vector<io_task<size_t>> reads;
702 : for (auto& buf : buffers)
703 : reads.push_back(stream.read_some(buf));
704 :
705 : auto result = co_await when_any(std::move(reads));
706 : if (result.index() == 1)
707 : {
708 : auto [idx, n] = std::get<1>(result);
709 : }
710 : }
711 : @endcode
712 :
713 : @see IoAwaitableRange, when_any
714 : */
715 : template<IoAwaitableRange R>
716 : requires detail::is_io_result_v<
717 : awaitable_result_t<std::ranges::range_value_t<R>>>
718 : && (!std::is_same_v<
719 : detail::io_result_payload_t<
720 : awaitable_result_t<std::ranges::range_value_t<R>>>,
721 : std::tuple<>>)
722 14 : [[nodiscard]] auto when_any(R&& awaitables)
723 : -> task<std::variant<std::error_code,
724 : std::pair<std::size_t,
725 : detail::io_result_payload_t<
726 : awaitable_result_t<std::ranges::range_value_t<R>>>>>>
727 : {
728 : using Awaitable = std::ranges::range_value_t<R>;
729 : using PayloadT = detail::io_result_payload_t<
730 : awaitable_result_t<Awaitable>>;
731 : using result_type = std::variant<std::error_code,
732 : std::pair<std::size_t, PayloadT>>;
733 : using OwnedRange = std::remove_cvref_t<R>;
734 :
735 : auto count = std::ranges::size(awaitables);
736 : if(count == 0)
737 : throw std::invalid_argument("when_any requires at least one awaitable");
738 :
739 : OwnedRange owned_awaitables = std::forward<R>(awaitables);
740 :
741 : detail::when_any_io_homogeneous_state<PayloadT> state(count);
742 :
743 : co_await detail::when_any_io_homogeneous_launcher<OwnedRange>(
744 : &owned_awaitables, &state);
745 :
746 : // Winner found
747 : if(state.core_.has_winner_.load(std::memory_order_acquire))
748 : {
749 : if(state.core_.winner_exception_)
750 : std::rethrow_exception(state.core_.winner_exception_);
751 : co_return result_type{std::in_place_index<1>,
752 : std::pair{state.core_.winner_index_, std::move(*state.result_)}};
753 : }
754 :
755 : // No winner — report the recorded failure
756 : if(state.last_exception_)
757 : std::rethrow_exception(state.last_exception_);
758 : co_return result_type{std::in_place_index<0>, state.last_error_};
759 28 : }
760 :
761 : /** Race a range of void io_result-returning awaitables.
762 :
763 : Only a child returning !ec can win. Returns the winner's index
764 : at variant index 1, or error_code at index 0 on all-fail.
765 :
766 : @par Await-effects
767 :
768 : Takes ownership of the range, creates one wrapper coroutine per
769 : element, then posts every wrapper to the caller's executor. All
770 : children therefore run concurrently, each awaited with the caller's
771 : executor and frame allocator and with a stop token owned by this
772 : operation.
773 :
774 : Awaiting an empty range throws `std::invalid_argument` before any
775 : child is started.
776 :
777 : The first child to await-return a zero `ec` claims the win. Claiming
778 : the win requests stop on the operation's own stop token, which every
779 : sibling observes through the stop token it was awaited with. A child
780 : that await-returns a non-zero `ec`, or that exits via an exception,
781 : does not claim the win and does not request stop. The operation keeps
782 : waiting for a success. A stop request on the caller's stop token is
783 : also forwarded to every child.
784 :
785 : The await completes only after every child has finished, regardless of
786 : whether a win was claimed.
787 :
788 : @par Await-returns
789 : An object of type `std::variant<std::error_code, std::size_t>`.
790 :
791 : @li Index 1 holds the winner's position in the input range. The
792 : children have no payloads, so nothing else is reported.
793 : @li Index 0 holds a non-zero `error_code` when no child won, that is,
794 : when every child failed. It is the `ec` of one of the failed
795 : children; which one is unspecified.
796 :
797 : If no child won and the failure selected for reporting is an
798 : exception rather than an `ec`, that exception is rethrown instead of
799 : await-returning. The choice of child is unspecified.
800 :
801 : @par Await-postcondition
802 : Every child has finished. The result holds index 1 if at least one
803 : child await-returned a zero `ec`. Otherwise the result holds index 0,
804 : or a failed child's exception is rethrown.
805 :
806 : @par Remarks
807 : Supports _IoAwaitable cancellation_. A canceled child await-returns a
808 : non-zero `ec` and so cannot win; if no child has already succeeded,
809 : the result settles at index 0.
810 :
811 : @par Thread Safety
812 : The returned task must be awaited from a single execution context.
813 : Child awaitables execute concurrently but complete through the caller's
814 : executor.
815 :
816 : @param awaitables Range of io_result<>-returning awaitables (must
817 : not be empty).
818 :
819 : @return A task yielding variant<error_code, size_t> where index 0
820 : is failure and index 1 carries the winner's index.
821 :
822 : @throws std::invalid_argument if range is empty.
823 :
824 : @par Exception Safety
825 : If all children fail and the reported failure is an exception,
826 : that child's exception is rethrown (which child is unspecified).
827 :
828 : @par Example
829 : @code
830 : task<void> example()
831 : {
832 : std::vector<io_task<>> jobs;
833 : jobs.push_back(background_work_a());
834 : jobs.push_back(background_work_b());
835 :
836 : auto result = co_await when_any(std::move(jobs));
837 : if (result.index() == 1)
838 : {
839 : auto winner = std::get<1>(result);
840 : }
841 : }
842 : @endcode
843 :
844 : @see IoAwaitableRange, when_any
845 : */
846 : template<IoAwaitableRange R>
847 : requires detail::is_io_result_v<
848 : awaitable_result_t<std::ranges::range_value_t<R>>>
849 : && std::is_same_v<
850 : detail::io_result_payload_t<
851 : awaitable_result_t<std::ranges::range_value_t<R>>>,
852 : std::tuple<>>
853 3 : [[nodiscard]] auto when_any(R&& awaitables)
854 : -> task<std::variant<std::error_code, std::size_t>>
855 : {
856 : using OwnedRange = std::remove_cvref_t<R>;
857 : using result_type = std::variant<std::error_code, std::size_t>;
858 :
859 : auto count = std::ranges::size(awaitables);
860 : if(count == 0)
861 : throw std::invalid_argument("when_any requires at least one awaitable");
862 :
863 : OwnedRange owned_awaitables = std::forward<R>(awaitables);
864 :
865 : detail::when_any_io_homogeneous_state<std::tuple<>> state(count);
866 :
867 : co_await detail::when_any_io_homogeneous_launcher<OwnedRange>(
868 : &owned_awaitables, &state);
869 :
870 : // Winner found
871 : if(state.core_.has_winner_.load(std::memory_order_acquire))
872 : {
873 : if(state.core_.winner_exception_)
874 : std::rethrow_exception(state.core_.winner_exception_);
875 : co_return result_type{std::in_place_index<1>,
876 : state.core_.winner_index_};
877 : }
878 :
879 : // No winner — report the recorded failure
880 : if(state.last_exception_)
881 : std::rethrow_exception(state.last_exception_);
882 : co_return result_type{std::in_place_index<0>, state.last_error_};
883 6 : }
884 :
885 : /** Race io_result-returning awaitables, selecting the first success.
886 :
887 : Overload selected when all children return io_result<Ts...>.
888 : Only a child returning !ec can win. Errors and exceptions do
889 : not claim winner status.
890 :
891 : @par Await-effects
892 :
893 : Creates and posts one wrapper coroutine per argument to the caller's
894 : executor, in argument order. All children therefore run concurrently,
895 : each awaited with the caller's executor and frame allocator and with
896 : a stop token owned by this operation. The overload requires at least
897 : one awaitable, so there is no empty case.
898 :
899 : The first child to await-return a zero `ec` claims the win. Claiming
900 : the win requests stop on the operation's own stop token, which every
901 : sibling observes through the stop token it was awaited with. A child
902 : that await-returns a non-zero `ec`, or that exits via an exception,
903 : does not claim the win and does not request stop. The operation keeps
904 : waiting for a success. A stop request on the caller's stop token is
905 : also forwarded to every child.
906 :
907 : The await completes only after every child has finished, regardless of
908 : whether a win was claimed.
909 :
910 : @par Await-returns
911 : An object of type `std::variant<std::error_code, P1, ..., Pn>`, where
912 : `Pi` is the payload of the i-th child's `io_result`.
913 :
914 : @li Index i+1 identifies the i-th argument as the winner and holds
915 : its payload.
916 : @li Index 0 holds a non-zero `error_code` when no child won, that is,
917 : when every child failed. It is the `ec` of one of the failed
918 : children; which one is unspecified.
919 :
920 : A child that succeeds after the win has already been claimed
921 : contributes nothing: its payload is discarded.
922 :
923 : If no child won and the failure selected for reporting is an
924 : exception rather than an `ec`, that exception is rethrown instead of
925 : await-returning. The choice of child is unspecified.
926 :
927 : @par Await-postcondition
928 : Every child has finished. If at least one child await-returned a zero
929 : `ec`, the result holds the index of the winning child, unless producing
930 : the winner's payload threw, in which case that exception is rethrown.
931 : Otherwise the result holds index 0, or a failed child's exception is
932 : rethrown.
933 :
934 : @par Remarks
935 : Supports _IoAwaitable cancellation_. A canceled child await-returns a
936 : non-zero `ec` and so cannot win; if no child has already succeeded,
937 : the result settles at index 0.
938 :
939 : @par Thread Safety
940 : The returned task must be awaited from a single execution context.
941 : Child awaitables execute concurrently but complete through the caller's
942 : executor.
943 :
944 : @param as The awaitables to race. Each must satisfy @ref
945 : IoAwaitable and is consumed (moved-from) when `when_any`
946 : is awaited.
947 :
948 : @return A task yielding variant<error_code, R1, ..., Rn> where
949 : index 0 is the failure/no-winner case and index i+1
950 : identifies the winning child. On all-fail, index 0 holds
951 : an error_code from one of the failed children (unspecified
952 : which; no priority between errors and exceptions).
953 :
954 : @par Exception Safety
955 : The winner's exception is rethrown if extracting or constructing
956 : the winning payload throws. In that case a winner was found, but
957 : its result could not be produced. If all children fail and the
958 : reported failure is an exception, that child's exception is
959 : rethrown (which child is unspecified).
960 :
961 : @note A failing child does not cancel its siblings; `when_any`
962 : waits for a success or for every child to finish. To make a
963 : benign error (e.g. @c cond::canceled) count as a win, wrap
964 : the child to translate the error into success. See the
965 : Concurrent Composition tutorial.
966 : */
967 : template<IoAwaitable... As>
968 : requires (sizeof...(As) > 0)
969 : && detail::all_io_result_awaitables<As...>
970 23 : [[nodiscard]] auto when_any(As... as)
971 : -> task<std::variant<
972 : std::error_code,
973 : detail::io_result_payload_t<awaitable_result_t<As>>...>>
974 : {
975 : using result_type = std::variant<
976 : std::error_code,
977 : detail::io_result_payload_t<awaitable_result_t<As>>...>;
978 :
979 : detail::when_any_io_state<
980 : detail::io_result_payload_t<awaitable_result_t<As>>...> state;
981 : std::tuple<As...> awaitable_tuple(std::move(as)...);
982 :
983 : co_await detail::when_any_io_launcher<As...>(
984 : &awaitable_tuple, &state);
985 :
986 : // Winner found: return their result
987 : if(state.result_.has_value())
988 : co_return std::move(*state.result_);
989 :
990 : // Winner claimed but payload construction failed
991 : if(state.core_.winner_exception_)
992 : std::rethrow_exception(state.core_.winner_exception_);
993 :
994 : // No winner — report the recorded failure
995 : if(state.last_exception_)
996 : std::rethrow_exception(state.last_exception_);
997 : co_return result_type{std::in_place_index<0>, state.last_error_};
998 46 : }
999 :
1000 : } // namespace capy
1001 : } // namespace boost
1002 :
1003 : #endif
|