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
21 changes: 16 additions & 5 deletions include/exec/sequence_senders.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -264,10 +264,15 @@ namespace experimental::execution
template <class _Env, class _Data>
struct __child_env_fn
{
auto operator()(_Env __env, _Data const & __data) const noexcept
-> STDEXEC::__join_env_t<_Data const &, _Env>
static_assert(STDEXEC::__nothrow_move_constructible<_Data>);

template <class _DataFwd>
auto operator()(_Env __env, _DataFwd&& __data) const
noexcept(STDEXEC::__nothrow_constructible_from<_Data, _DataFwd&&>)
-> STDEXEC::__join_env_t<_Data, _Env>
{
return STDEXEC::__env::__join(__data, static_cast<_Env&&>(__env));
return STDEXEC::__env::__join(_Data(static_cast<_DataFwd&&>(__data)),
static_cast<_Env&&>(__env));
}
};
};
Expand Down Expand Up @@ -1010,8 +1015,14 @@ namespace experimental::execution
using __rcvr_t = __adaptor_rcvr<_Receiver, __child_env_t>;
using __result_t = STDEXEC::__call_result_t<subscribe_t, __child_t, __rcvr_t>;
__check_operation_state<__result_t>();
using __xform_t = __adaptor_child_env_fn_t<__tag_t, env_of_t<_Receiver>, __data_t>;
using __data_fwd_t = decltype(STDEXEC::__forward_like<__tfx_seq_t>(
__declval<__data_t&>()));
constexpr bool __nothrow_subscribe = __nothrow_callable<subscribe_t, __child_t, __rcvr_t>;
return __declfn<__result_t, __nothrow_subscribe && __nothrow_tfx_seq>();
constexpr bool __nothrow_child_env =
__nothrow_callable<__xform_t, env_of_t<_Receiver>, __data_fwd_t>;
return __declfn<__result_t,
__nothrow_subscribe && __nothrow_child_env && __nothrow_tfx_seq>();
}
else if constexpr (__subscribable_with_static_member<__tfx_seq_t, _Receiver>)
{
Expand Down Expand Up @@ -1084,7 +1095,7 @@ namespace experimental::execution
__adaptor_rcvr<_Receiver, __child_env_t>{
static_cast<_Receiver&&>(__rcvr),
__xform_t{}(static_cast<decltype(__env)&&>(__env),
STDEXEC::__forward_like<decltype(__data)>(__data))});
STDEXEC::__forward_like<__tfx_seq_t>(__data))});
}
else if constexpr (__subscribable_with_static_member<__tfx_seq_t, _Receiver>)
{ // NOLINT(bugprone-branch-clone)
Expand Down
152 changes: 152 additions & 0 deletions test/exec/sequence/test_write_env_sequence.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@
#include <test_common/senders.hpp>
#include <test_common/type_helpers.hpp>

#include <optional>
#include <utility>

namespace
Expand Down Expand Up @@ -150,4 +151,155 @@ namespace
| exec::ignore_all_values());
CHECK(value == 42);
}

// Lifetime tests for issue #2305: `write_env` over a sequence sender injects
// its data into the environment the child is subscribed with, and the
// resulting operation state must own that data rather than reference the
// (long-dead) sender.

struct read_env_query : STDEXEC::__query<read_env_query>
{
static consteval auto query(STDEXEC::forwarding_query_t) noexcept -> bool
{
return true;
}
};

inline constexpr int poisoned_value = -1;

struct injected_env_data
{
int value;

explicit injected_env_data(int value) noexcept
: value(value)
{}

injected_env_data(injected_env_data const &) = default;
injected_env_data(injected_env_data&&) noexcept = default;

// Poison the value on destruction so that a read through a dangling
// reference after the sender is gone fails deterministically.
~injected_env_data()
{
value = poisoned_value;
}

auto query(read_env_query) const noexcept -> int
{
return value;
}
};

struct move_only_env_data
{
int value;

explicit move_only_env_data(int value) noexcept
: value(value)
{}

move_only_env_data(move_only_env_data&&) noexcept = default;
move_only_env_data(move_only_env_data const &) = delete;

~move_only_env_data()
{
value = poisoned_value;
}

auto query(read_env_query) const noexcept -> int
{
return value;
}
};

// A sequence sender that reads an injected query out of its receiver's
// environment when started.
struct env_reading_sequence
{
using sender_concept = exec::sequence_sender_tag;
using item_types = exec::item_types<decltype(STDEXEC::just(int{}))>;
using completion_signatures =
STDEXEC::completion_signatures<STDEXEC::set_value_t(), STDEXEC::set_stopped_t()>;

int* observed_;

template <class Rcvr>
struct op
{
using operation_state_concept = STDEXEC::operation_state_t;
Rcvr rcvr_;
int* observed_;

void start() noexcept
{
*observed_ = read_env_query{}(STDEXEC::get_env(rcvr_));
STDEXEC::set_value(static_cast<Rcvr&&>(rcvr_));
}
};

template <STDEXEC::receiver Rcvr>
auto subscribe(Rcvr rcvr) const -> op<Rcvr>
{
return op<Rcvr>{static_cast<Rcvr&&>(rcvr), observed_};
}
};

struct test_sequence_rcvr
{
using receiver_concept = STDEXEC::receiver_tag;

template <class Item>
auto set_next(Item&& item) noexcept
{
return static_cast<Item&&>(item);
}

void set_value() noexcept {}

void set_stopped() noexcept {}

template <class Error>
void set_error(Error&&) noexcept
{}

auto get_env() const noexcept
{
return STDEXEC::prop(STDEXEC::get_stop_token, STDEXEC::inplace_stop_token{});
}
};

TEST_CASE("write_env's env data outlives a temporary sequence sender", "[sequence][write_env]")
{
int observed = 0;
auto op = exec::subscribe(STDEXEC::write_env(env_reading_sequence{&observed},
injected_env_data{42}),
test_sequence_rcvr{});
STDEXEC::start(op);
CHECK(observed == 42);
}

TEST_CASE("write_env with a move-only env over a sequence sender", "[sequence][write_env]")
{
int observed = 0;
auto op = exec::subscribe(STDEXEC::write_env(env_reading_sequence{&observed},
move_only_env_data{42}),
test_sequence_rcvr{});
STDEXEC::start(op);
CHECK(observed == 42);
}

TEST_CASE("write_env's env data outlives a named sequence sender", "[sequence][write_env]")
{
int observed = 0;
using sndr_t = decltype(STDEXEC::write_env(env_reading_sequence{}, injected_env_data{0}));
using op_t = exec::subscribe_result_t<sndr_t&, test_sequence_rcvr>;
std::optional<sndr_t> sndr;
sndr.emplace(STDEXEC::write_env(env_reading_sequence{&observed}, injected_env_data{42}));
std::optional<op_t> op;
op.emplace(exec::subscribe(*sndr, test_sequence_rcvr{}));
sndr.reset();
STDEXEC::start(*op);
CHECK(observed == 42);
}
} // namespace