Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions include/exec/reschedule.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@ namespace experimental::execution
{
struct _CANNOT_RESCHEDULE_
{};
using STDEXEC::_THE_CURRENT_EXECUTION_ENVIRONMENT_DOESNT_HAVE_A_SCHEDULER_;
using STDEXEC::_THE_CURRENT_EXECUTION_ENVIRONMENT_DOES_NOT_HAVE_A_START_SCHEDULER_;

namespace __resched
{
Expand All @@ -31,7 +31,7 @@ namespace experimental::execution
template <class _Env>
using __no_scheduler_error =
__mexception<_WHAT_(_CANNOT_RESCHEDULE_),
_WHY_(_THE_CURRENT_EXECUTION_ENVIRONMENT_DOESNT_HAVE_A_SCHEDULER_),
_WHY_(_THE_CURRENT_EXECUTION_ENVIRONMENT_DOES_NOT_HAVE_A_START_SCHEDULER_),
_WHERE_(_IN_ALGORITHM_, reschedule_t),
_WITH_ENVIRONMENT_(_Env)>;

Expand Down
2 changes: 1 addition & 1 deletion include/stdexec/__detail/__affine.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -109,7 +109,7 @@ namespace STDEXEC
// sender to be affine. Instead, return a type describing the problem.
return __not_a_sender< //
_WHAT_(_CANNOT_MAKE_SENDER_AFFINE_TO_THE_STARTING_SCHEDULER_),
_WHY_(_THE_CURRENT_EXECUTION_ENVIRONMENT_DOESNT_HAVE_A_SCHEDULER_),
_WHY_(_THE_CURRENT_EXECUTION_ENVIRONMENT_DOES_NOT_HAVE_A_START_SCHEDULER_),
_WHERE_(_IN_ALGORITHM_, affine_t),
_WITH_PRETTY_SENDER_<__cv_child_t>,
_WITH_ENVIRONMENT_(_Env)>{};
Expand Down
83 changes: 54 additions & 29 deletions include/stdexec/__detail/__continues_on.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -48,13 +48,20 @@ namespace STDEXEC
// [exec.continues.on]
namespace __trnsfr
{
template <class _Sexpr, class _Receiver>
template <class _Sender, class _Receiver>
struct __state_base
{
using __storage_t = __storage_for_t<__child_of<_Sexpr>, env_of_t<_Receiver>>;
using __env2_t = __secondary_env_t<_Sender, env_of_t<_Receiver>, set_value_t>;
using __storage_t = __storage_for_t<_Sender, env_of_t<_Receiver>>;

_Receiver __rcvr_;
__storage_t __data_;
constexpr __state_base(_Sender const & __child, _Receiver&& __rcvr) noexcept
: __rcvr_(static_cast<_Receiver&&>(__rcvr))
, __env2_(__mk_secondary_env_t<set_value_t>()(__child, get_env(__rcvr_)))
{}

_Receiver __rcvr_;
__env2_t const __env2_;
__storage_t __data_;
};

// This receiver is to be completed on the execution context associated with
Expand All @@ -63,10 +70,12 @@ namespace STDEXEC
// receiver completes, it can read the completion out of the operation state
// and forward it to the output receiver after transitioning to the
// scheduler's context.
template <class _Sexpr, class _Receiver>
template <class _Sender, class _Receiver>
struct __receiver2
{
using receiver_concept = receiver_tag;
using __env2_t = __state_base<_Sender, _Receiver>::__env2_t;
using __sch_env_t = __join_env_t<__env2_t const &, env_of_t<_Receiver>>;

constexpr void set_value() noexcept
{
Expand All @@ -86,24 +95,24 @@ namespace STDEXEC
}

[[nodiscard]]
constexpr auto get_env() const noexcept -> env_of_t<_Receiver>
constexpr auto get_env() const noexcept -> __sch_env_t
{
return STDEXEC::get_env(__state_->__rcvr_);
return __env::__join(__state_->__env2_, STDEXEC::get_env(__state_->__rcvr_));
}

__state_base<_Sexpr, _Receiver>* __state_;
__state_base<_Sender, _Receiver>* __state_;
};

template <class _Scheduler, class _Sexpr, class _Receiver>
struct __state : __state_base<_Sexpr, _Receiver>
template <class _Scheduler, class _Sender, class _Receiver>
struct __state : __state_base<_Sender, _Receiver>
{
using __receiver2_t = __receiver2<_Sexpr, _Receiver>;
using __receiver2_t = __receiver2<_Sender, _Receiver>;
using __schedule_sender_t = schedule_result_t<_Scheduler&>;

constexpr explicit __state(_Scheduler __sched, _Receiver&& __rcvr)
constexpr explicit __state(_Scheduler __sched, _Sender const & __child, _Receiver&& __rcvr)
noexcept(__nothrow_callable<schedule_t, _Scheduler&>
&& __nothrow_connectable<__schedule_sender_t, __receiver2_t>)
: __state::__state_base{static_cast<_Receiver&&>(__rcvr)}
: __state::__state_base{__child, static_cast<_Receiver&&>(__rcvr)}
, __state2_(connect(schedule(__sched), __receiver2_t{this}))
{}
STDEXEC_IMMOVABLE(__state);
Expand All @@ -116,6 +125,18 @@ namespace STDEXEC
struct __attrs
{
private:
template <class _Env>
using __env2_t = __secondary_env_t<_Sender, _Env, set_value_t>;
template <class _Env>
using __sch_env_t = __join_env_t<__env2_t<_Env>, _Env>;

template <class _Env>
constexpr auto __mk_sch_env(_Env&& __env) const noexcept -> __sch_env_t<_Env>
{
return __env::__join(__mk_secondary_env_t<set_value_t>()(__sndr_, __env),
static_cast<_Env&&>(__env));
}
Comment on lines +128 to +138

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Noting that the tag is hardcoded here to set_value_t, whereas the call sites are templated on _SetTag? Could this lead to issues e.g. when the predecessor completes with error/stop?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think this issue is still open and may be relevant.

The __complete lambda stores all completions with their tag, along with a comment that states that the goal is to forward them from within the scheduler's context.

https://github.com/ericniebler/stdexec/blob/233768b59e2d92c65c6652a6d18864db15e8f1ac/include/stdexec/__detail/__continues_on.hpp#L362-L369

This is indeed what __receiver2 does:

https://github.com/ericniebler/stdexec/blob/233768b59e2d92c65c6652a6d18864db15e8f1ac/include/stdexec/__detail/__continues_on.hpp#L79-L82

If the predecessor completes with an error, we can thus expect that error to be re-emitted from the scheduler's context, and it seems that's a case that the queries don't capture yet.

Incidentally, the first sentence of the docstring of __receiver2 appears too strong because its set_error and set_stopped would be invoked from the predecessor's context.


//! @brief Returns `true` when:
//! - _SetTag is set_error_t, and
//! - _Sender has value completions, and
Expand All @@ -138,13 +159,13 @@ namespace STDEXEC
return false;
}

_Scheduler __sch_;
env_of_t<_Sender> __attrs_;
_Scheduler __sch_;
_Sender const & __sndr_;

public:
constexpr explicit __attrs(_Scheduler __sch, env_of_t<_Sender> __attrs) noexcept
constexpr explicit __attrs(_Scheduler __sch, _Sender const & __sndr) noexcept
: __sch_(static_cast<_Scheduler&&>(__sch))
, __attrs_(static_cast<env_of_t<_Sender>&&>(__attrs))
, __sndr_(__sndr)
{}

//! @brief Queries the completion scheduler for a given @c _SetTag.
Expand All @@ -171,22 +192,22 @@ namespace STDEXEC
[[nodiscard]]
constexpr auto
query(get_completion_scheduler_t<_SetTag>, _Env const &... __env) const noexcept
-> __call_result_t<get_completion_scheduler_t<_SetTag>, _Scheduler, __fwd_env_t<_Env>...>
-> __call_result_t<get_completion_scheduler_t<_SetTag>, _Scheduler, __sch_env_t<_Env>...>
{
return get_completion_scheduler<_SetTag>(__sch_, __fwd_env(__env)...);
return get_completion_scheduler<_SetTag>(__sch_, __mk_sch_env(__env)...);
}

//! @overload
template <class _SetTag, class... _Env>
requires __never_sends<_SetTag, schedule_result_t<_Scheduler>, __fwd_env_t<_Env>...>
requires __never_sends<_SetTag, schedule_result_t<_Scheduler>, __sch_env_t<_Env>...>
[[nodiscard]]
constexpr auto
query(get_completion_scheduler_t<_SetTag>, _Env const &... __env) const noexcept
-> __call_result_t<get_completion_scheduler_t<_SetTag>,
env_of_t<_Sender>,
__fwd_env_t<_Env>...>
{
return get_completion_scheduler<_SetTag>(__attrs_, __fwd_env(__env)...);
return get_completion_scheduler<_SetTag>(get_env(__sndr_), __fwd_env(__env)...);
}

//! @brief Queries the completion domain for a given @c _SetTag.
Expand All @@ -209,7 +230,7 @@ namespace STDEXEC
[[nodiscard]]
constexpr auto
query(get_completion_domain_t<_SetTag>, _Env const &...) const noexcept -> __unless_one_of_t<
__completion_domain_of_t<_SetTag, schedule_result_t<_Scheduler>, __fwd_env_t<_Env>...>,
__completion_domain_of_t<_SetTag, schedule_result_t<_Scheduler>, __env2_t<_Env>...>,
indeterminate_domain<>>
{
return {};
Expand All @@ -223,7 +244,7 @@ namespace STDEXEC
query(get_completion_domain_t<_SetTag>, _Env const &...) const noexcept -> __unless_one_of_t<
__common_domain_t<
__completion_domain_of_t<_SetTag, _Sender, __fwd_env_t<_Env>...>,
__completion_domain_of_t<_SetTag, schedule_result_t<_Scheduler>, __fwd_env_t<_Env>...>>,
__completion_domain_of_t<_SetTag, schedule_result_t<_Scheduler>, __env2_t<_Env>...>>,
indeterminate_domain<>>
{
return {};
Expand All @@ -237,7 +258,7 @@ namespace STDEXEC
query(get_completion_domain_t<_SetTag>, _Env const &...) const noexcept -> __unless_one_of_t<
__common_domain_t<
__completion_domain_of_t<_SetTag, _Sender, __fwd_env_t<_Env>...>,
__completion_domain_of_t<_SetTag, schedule_result_t<_Scheduler>, __fwd_env_t<_Env>...>,
__completion_domain_of_t<_SetTag, schedule_result_t<_Scheduler>, __env2_t<_Env>...>,
__completion_domain_of_t<set_value_t, _Sender, __fwd_env_t<_Env>...>>,
indeterminate_domain<>>
{
Expand All @@ -255,7 +276,7 @@ namespace STDEXEC
{
using _SchSender = schedule_result_t<_Scheduler>;
constexpr auto cb_sched =
STDEXEC::__get_completion_behavior<_Tag, _SchSender, __fwd_env_t<_Env>...>();
STDEXEC::__get_completion_behavior<_Tag, _SchSender, __env2_t<_Env>...>();
constexpr auto cb_sndr =
STDEXEC::__get_completion_behavior<_Tag, _Sender, __fwd_env_t<_Env>...>();
return cb_sched | cb_sndr;
Expand All @@ -271,7 +292,7 @@ namespace STDEXEC
noexcept(__nothrow_queryable_with<env_of_t<_Sender>, _Query, _Args...>)
-> __query_result_t<env_of_t<_Sender>, _Query, _Args...>
{
return __attrs_.query(_Query(), static_cast<_Args&&>(__args)...);
return _Query()(STDEXEC::get_env(__sndr_), static_cast<_Args&&>(__args)...);
}
};

Expand Down Expand Up @@ -303,15 +324,16 @@ namespace STDEXEC
}

template <class _Sender, class _Receiver>
using __state_for_t = __state<__decay_t<__data_of<_Sender>>, _Sender, _Receiver>;
using __state_for_t =
__state<__decay_t<__data_of<_Sender>>, __decay_t<__child_of<_Sender>>, _Receiver>;

public:
static constexpr auto __get_attrs =
[]<class _Scheduler, class _Child>(__ignore,
_Scheduler const & __data,
_Child const & __child) noexcept
{
return __attrs<_Scheduler, _Child>{__data, STDEXEC::get_env(__child)};
return __attrs<_Scheduler, _Child>{__data, __child};
};

template <class _Sender, class... _Env>
Expand All @@ -332,12 +354,15 @@ namespace STDEXEC
[]<class _Sender, class _Receiver>(_Sender&& __sndr, _Receiver&& __rcvr) noexcept(
__nothrow_constructible_from<__state_for_t<_Sender, _Receiver>,
__data_of<_Sender>&,
__child_of<_Sender>&,
_Receiver>) -> __state_for_t<_Sender, _Receiver>
requires sender_in<__child_of<_Sender>, __fwd_env_t<env_of_t<_Receiver>>>
{
static_assert(__sender_for<_Sender, continues_on_t>);
auto& [__tag, __sched, __child] = __sndr;
return __state_for_t<_Sender, _Receiver>{__sched, static_cast<_Receiver&&>(__rcvr)};
return __state_for_t<_Sender, _Receiver>{__sched,
__child,
static_cast<_Receiver&&>(__rcvr)};
};

static constexpr auto __complete =
Expand Down
2 changes: 1 addition & 1 deletion include/stdexec/__detail/__counting_scopes.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -116,7 +116,7 @@ namespace STDEXEC
{
return STDEXEC::__throw_compile_time_error<
_WHAT_(_JOINING_A_COUNTING_SCOPE_NEEDS_A_SCHEDULER_IN_THE_ENVIRONMENT_),
_WHY_(_THE_CURRENT_EXECUTION_ENVIRONMENT_DOESNT_HAVE_A_SCHEDULER_),
_WHY_(_THE_CURRENT_EXECUTION_ENVIRONMENT_DOES_NOT_HAVE_A_START_SCHEDULER_),
_WHERE_(STDEXEC::_IN_ALGORITHM_, __scope_join_t),
_WITH_PRETTY_SENDER_<_Sender>,
_WITH_ENVIRONMENT_(_Env)>();
Expand Down
5 changes: 4 additions & 1 deletion include/stdexec/__detail/__diagnostics.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -151,7 +151,10 @@ namespace STDEXEC
struct _CONCEPT_CHECK_FAILURE_
{};

struct _THE_CURRENT_EXECUTION_ENVIRONMENT_DOESNT_HAVE_A_SCHEDULER_
struct _THE_CURRENT_EXECUTION_ENVIRONMENT_DOES_NOT_HAVE_A_START_SCHEDULER_
{};

struct _THE_PREDECESSOR_SENDER_DOES_NOT_KNOW_THE_SCHEDULER_ON_WHICH_IT_WILL_COMPLETE_
{};

template <class _Sender>
Expand Down
7 changes: 3 additions & 4 deletions include/stdexec/__detail/__env.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -238,8 +238,7 @@ namespace STDEXEC
STDEXEC_ATTRIBUTE(nodiscard, always_inline, host, device)
constexpr auto operator()(env<_Envs...> const &__env) const noexcept -> decltype(auto)
{
// count of elements that includes the first env that supports the query
// and all subsequent envs
// compute the index of the first env that supports the query:
STDEXEC_CONSTEXPR_LOCAL auto __index =
sizeof...(_Envs) - __mcall<__mfind_if<__q1<__has_query_t>, __msize>, _Envs...>::value;
if constexpr (__index < sizeof...(_Envs))
Expand All @@ -265,8 +264,8 @@ namespace STDEXEC
noexcept(__nothrow_queryable_with<__1st_env_t<_Query, _Args...>, _Query, _Args...>)
-> __query_result_t<__1st_env_t<_Query, _Args...>, _Query, _Args...>
{
auto const &__env = __detail::__get_1st_env<_Query, _Args...>()(*this);
return __env.query(_Query(), static_cast<_Args &&>(__args)...);
constexpr auto __get_env = __detail::__get_1st_env<_Query, _Args...>();
return __get_env(*this).query(_Query(), static_cast<_Args &&>(__args)...);
}
};

Expand Down
Loading
Loading