diff --git a/include/exec/timed_scheduler.hpp b/include/exec/timed_scheduler.hpp index 5fac8a18a..5b2d4de3a 100644 --- a/include/exec/timed_scheduler.hpp +++ b/include/exec/timed_scheduler.hpp @@ -20,6 +20,79 @@ #include +namespace experimental::execution::__timed_scheduler_fallback +{ + using namespace STDEXEC; + + struct __tag + {}; + + template + concept __completion_scheduler_matches = + __callable, env_of_t<_NativeSender>> + && __decays_to< + __call_result_t, env_of_t<_NativeSender>>, + __decay_t<_Scheduler>>; + + template + struct __data + { + _Scheduler __sched_; + }; + + template + struct __attrs : STDEXEC::__sync_attrs<_Sender> + { + using __base_t = STDEXEC::__sync_attrs<_Sender>; + using __base_t::query; + + constexpr __attrs(_Scheduler __sched, _Sender const &__sndr) noexcept + : __base_t{__sndr} + , __sched_{static_cast<_Scheduler &&>(__sched)} + {} + + [[nodiscard]] + constexpr auto + query(STDEXEC::get_completion_scheduler_t) const noexcept -> _Scheduler + requires __completion_scheduler_matches<_NativeSender, _Scheduler> + { + return __sched_; + } + + _Scheduler __sched_; + }; +} // namespace experimental::execution::__timed_scheduler_fallback + +namespace STDEXEC +{ + template <> + struct __sexpr_impl<::experimental::execution::__timed_scheduler_fallback::__tag> + : __sexpr_defaults + { + static constexpr auto __get_attrs = + []( + __ignore, + ::experimental::execution::__timed_scheduler_fallback::__data<_Scheduler, + _NativeSender> const &__data, + _Sender const &__sndr) noexcept + { + return ::experimental::execution::__timed_scheduler_fallback::__attrs<_Scheduler, + _NativeSender, + _Sender>{ + __data.__sched_, + __sndr}; + }; + + template + static consteval auto __get_completion_signatures() + { + static_assert( + __sender_for<_Sender, ::experimental::execution::__timed_scheduler_fallback::__tag>); + return STDEXEC::get_completion_signatures<__child_of<_Sender>, _Env...>(); + } + }; +} // namespace STDEXEC + namespace experimental::execution { namespace __now @@ -167,13 +240,17 @@ namespace experimental::execution auto operator()(_Scheduler &&__sched, const duration_of_t<_Scheduler> &__duration) const noexcept { - // TODO get_completion_scheduler - return let_value( - just(), - [__sched, __duration]() noexcept( - __nothrow_callable> - &&__nothrow_callable) - { return schedule_at(__sched, now(__sched) + __duration); }); + using __native_sender_t = + __call_result_t<__schedule_at_base_t, _Scheduler, time_point_of_t<_Scheduler> const &>; + + return __make_sexpr<__timed_scheduler_fallback::__tag>( + __timed_scheduler_fallback::__data, __native_sender_t>{ + __sched}, + let_value(just(), + [__sched, __duration]() noexcept( + __nothrow_callable> + &&__nothrow_callable) + { return schedule_at(__sched, now(__sched) + __duration); })); } }; } // namespace __schedule_after @@ -244,11 +321,16 @@ namespace experimental::execution auto operator()(_Scheduler &&__sched, const time_point_of_t<_Scheduler> &__time_point) const noexcept(noexcept(schedule_after(__sched, __time_point - now(__sched)))) { - // TODO get_completion_scheduler - return let_value(just(), - [__sched, __time_point]() noexcept( - noexcept(schedule_after(__sched, __time_point - now(__sched)))) - { return schedule_after(__sched, __time_point - now(__sched)); }); + using __native_sender_t = + __call_result_t<__schedule_after_base_t, _Scheduler, duration_of_t<_Scheduler> const &>; + + return __make_sexpr<__timed_scheduler_fallback::__tag>( + __timed_scheduler_fallback::__data, __native_sender_t>{ + __sched}, + let_value(just(), + [__sched, __time_point]() noexcept( + noexcept(schedule_after(__sched, __time_point - now(__sched)))) + { return schedule_after(__sched, __time_point - now(__sched)); })); } }; } // namespace __schedule_at diff --git a/test/exec/test_timed_thread_scheduler.cpp b/test/exec/test_timed_thread_scheduler.cpp index ad4d13b74..cd2b9ca58 100644 --- a/test/exec/test_timed_thread_scheduler.cpp +++ b/test/exec/test_timed_thread_scheduler.cpp @@ -31,6 +31,177 @@ #else namespace { + struct schedule_after_only_scheduler + { + using scheduler_concept = STDEXEC::scheduler_tag; + using time_point = std::chrono::steady_clock::time_point; + using duration = time_point::duration; + + struct sender; + + constexpr auto + operator==(schedule_after_only_scheduler const &) const noexcept -> bool = default; + + [[nodiscard]] + auto now() const noexcept -> time_point + { + return std::chrono::steady_clock::now(); + } + + [[nodiscard]] + auto schedule() const noexcept -> sender; + + [[nodiscard]] + auto schedule_after(duration) const noexcept -> sender; + }; + + struct schedule_after_only_scheduler::sender + { + using sender_concept = STDEXEC::sender_tag; + using completion_signatures = STDEXEC::completion_signatures; + + struct attrs + { + [[nodiscard]] + auto query(STDEXEC::get_completion_scheduler_t) const noexcept + -> schedule_after_only_scheduler + { + return {}; + } + }; + + template + struct operation + { + void start() & noexcept + { + STDEXEC::set_value(static_cast(receiver_)); + } + + Receiver receiver_; + }; + + template + auto connect(Receiver receiver) const -> operation + { + return {static_cast(receiver)}; + } + + [[nodiscard]] + constexpr auto get_env() const noexcept -> attrs + { + return {}; + } + }; + + auto schedule_after_only_scheduler::schedule() const noexcept -> sender + { + return {}; + } + + auto schedule_after_only_scheduler::schedule_after(duration) const noexcept -> sender + { + return {}; + } + + template + struct sender_with_completion_scheduler + { + using sender_concept = STDEXEC::sender_tag; + using completion_signatures = STDEXEC::completion_signatures; + + struct attrs + { + [[nodiscard]] + auto query(STDEXEC::get_completion_scheduler_t) const noexcept + -> CompletionScheduler + { + return {}; + } + }; + + template + struct operation + { + void start() & noexcept + { + STDEXEC::set_value(static_cast(receiver_)); + } + + Receiver receiver_; + }; + + template + auto connect(Receiver receiver) const -> operation + { + return {static_cast(receiver)}; + } + + [[nodiscard]] + constexpr auto get_env() const noexcept -> attrs + { + return {}; + } + }; + + struct schedule_after_delegating_scheduler + { + using scheduler_concept = STDEXEC::scheduler_tag; + using time_point = std::chrono::steady_clock::time_point; + using duration = time_point::duration; + using sender = sender_with_completion_scheduler; + + constexpr auto + operator==(schedule_after_delegating_scheduler const &) const noexcept -> bool = default; + + [[nodiscard]] + auto now() const noexcept -> time_point + { + return std::chrono::steady_clock::now(); + } + + [[nodiscard]] + auto schedule() const noexcept -> sender + { + return {}; + } + + [[nodiscard]] + auto schedule_after(duration) const noexcept -> sender + { + return {}; + } + }; + + struct schedule_at_delegating_scheduler + { + using scheduler_concept = STDEXEC::scheduler_tag; + using time_point = std::chrono::steady_clock::time_point; + using duration = time_point::duration; + using sender = sender_with_completion_scheduler; + + constexpr auto + operator==(schedule_at_delegating_scheduler const &) const noexcept -> bool = default; + + [[nodiscard]] + auto now() const noexcept -> time_point + { + return std::chrono::steady_clock::now(); + } + + [[nodiscard]] + auto schedule() const noexcept -> sender + { + return {}; + } + + [[nodiscard]] + auto schedule_at(time_point) const noexcept -> sender + { + return {}; + } + }; + TEST_CASE("timed_thread_scheduler - unused context", "[types][timed_thread_scheduler][schedulers]") { @@ -80,6 +251,51 @@ namespace CHECK(STDEXEC::sync_wait(exec::schedule_after(scheduler, duration))); } + TEST_CASE("timed scheduler fallbacks advertise completion schedulers", + "[timed_scheduler][completion_scheduler]") + { + static_assert(exec::timed_scheduler); + + exec::timed_thread_context context; + exec::timed_thread_scheduler timed_scheduler = context.get_scheduler(); + auto after_sender = exec::schedule_after(timed_scheduler, std::chrono::milliseconds(10)); + + CHECK(STDEXEC::get_completion_scheduler(STDEXEC::get_env(after_sender)) + == timed_scheduler); + + schedule_after_only_scheduler after_only_scheduler; + auto at_sender = exec::schedule_at(after_only_scheduler, exec::now(after_only_scheduler)); + + CHECK(STDEXEC::get_completion_scheduler(STDEXEC::get_env(at_sender)) + == after_only_scheduler); + CHECK(STDEXEC::sync_wait(std::move(at_sender))); + } + + TEST_CASE("timed scheduler fallbacks do not invent completion schedulers", + "[timed_scheduler][completion_scheduler]") + { + using completion_scheduler_query = STDEXEC::get_completion_scheduler_t; + + static_assert(exec::timed_scheduler); + using at_sender_t = decltype(exec::schedule_at( + std::declval(), + std::declval const &>())); + static_assert(!STDEXEC::__callable>); + + static_assert(exec::timed_scheduler); + using after_sender_t = decltype(exec::schedule_after( + std::declval(), + std::declval const &>())); + static_assert( + !STDEXEC::__callable>); + + schedule_after_delegating_scheduler after_scheduler; + CHECK(STDEXEC::sync_wait(exec::schedule_at(after_scheduler, exec::now(after_scheduler)))); + + schedule_at_delegating_scheduler at_scheduler; + CHECK(STDEXEC::sync_wait(exec::schedule_after(at_scheduler, std::chrono::milliseconds(0)))); + } + TEST_CASE("timed_thread_scheduler - when_any", "[timed_thread_scheduler][when_any]") { exec::timed_thread_context context;