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
44 changes: 41 additions & 3 deletions include/exec/fork_join.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -39,8 +39,24 @@ namespace experimental::execution
}
};

template <class Completions,
bool = STDEXEC::__nothrow_decay_copyable_results_t<Completions>::value>
struct _cache_completions
{
using type = Completions;
};

template <class Completions>
struct _cache_completions<Completions, false>
{
using type =
STDEXEC::__concat_completion_signatures_t<Completions, STDEXEC::__eptr_completion_t>;
};

// Account that decay-copying into the cache may throw.
template <class Completions>
using _variant_t = STDEXEC::__mapply_q<STDEXEC::__results_storage, Completions>;
using _variant_t = STDEXEC::__mapply_q<STDEXEC::__results_storage,
typename _cache_completions<Completions>::type>;

template <class Domain>
struct _env_t
Expand Down Expand Up @@ -281,16 +297,38 @@ namespace experimental::execution

struct fork_join_t
{
/// No closure given.
template <STDEXEC::sender Sndr>
STDEXEC_ATTRIBUTE(host, device)
constexpr auto operator()(Sndr&& sndr) const noexcept(STDEXEC::__nothrow_decay_copyable<Sndr>)
{
return static_cast<Sndr&&>(sndr);
}

/// Unary closure.
template <STDEXEC::sender Sndr, class Closure>
requires(!STDEXEC::sender<Closure>)
STDEXEC_ATTRIBUTE(host, device)
constexpr auto operator()(Sndr&& sndr, Closure&& clsr) const
Comment thread
romintomasetti marked this conversation as resolved.
noexcept(STDEXEC::__nothrow_callable<Closure, Sndr>)
{
return static_cast<Closure&&>(clsr)(static_cast<Sndr&&>(sndr));
}

/// One sender and multiple closures.
template <STDEXEC::sender Sndr, class... Closures>
requires(sizeof...(Closures) > 1)
STDEXEC_ATTRIBUTE(host, device)
constexpr auto operator()(Sndr&& sndr, Closures&&... closures) const //
-> STDEXEC::__well_formed_sender auto
constexpr auto operator()(Sndr&& sndr, Closures&&... closures) const
noexcept(STDEXEC::__nothrow_decay_copyable<Sndr, Closures...>)
-> STDEXEC::__well_formed_sender auto
{
return STDEXEC::__sexpr{fork_join_t(),
STDEXEC::__tuple{static_cast<Closures&&>(closures)...},
static_cast<Sndr&&>(sndr)};
}

/// One or more closures.
template <class... Closures>
requires((!STDEXEC::sender<Closures>) && ...)
STDEXEC_ATTRIBUTE(host, device)
Expand Down
2 changes: 1 addition & 1 deletion include/stdexec/__detail/__then.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -239,7 +239,7 @@ namespace STDEXEC
//! sender `then(sndr, std::move(__fun))`.
template <__movable_value _Fun>
STDEXEC_ATTRIBUTE(always_inline, host, device)
constexpr auto operator()(_Fun __fun) const
constexpr auto operator()(_Fun __fun) const noexcept(__nothrow_decay_copyable<_Fun>)
{
return __closure(*this, static_cast<_Fun&&>(__fun));
}
Expand Down
80 changes: 79 additions & 1 deletion test/exec/test_fork_join.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,30 @@ namespace
completion_signatures<set_value_t(), set_error_t(std::exception_ptr)>>);
}

TEST_CASE("fork_join coalesces empty and unary calls", "[adaptors][fork_join]")
{
/// Empty (no closure given).
STDEXEC::sender auto empty = exec::fork_join(STDEXEC::just());
using empty_t = decltype(empty);
STATIC_REQUIRE(std::same_as<empty_t, decltype(STDEXEC::just())>);
STATIC_REQUIRE(!exec::sender_for<empty_t, exec::fork_join_t>);
STATIC_REQUIRE(noexcept(exec::fork_join(STDEXEC::just())));

auto then = STDEXEC::then([]() noexcept {});

/// Unary closure.
STDEXEC::sender auto unary = exec::fork_join(STDEXEC::just(), then);
using unary_t = decltype(unary);
STATIC_REQUIRE(std::same_as<unary_t, decltype(STDEXEC::just() | then)>);
STATIC_REQUIRE(!exec::sender_for<unary_t, exec::fork_join_t>);
STATIC_REQUIRE(noexcept(exec::fork_join(STDEXEC::just(), then)));

/// Multiple closures.
STDEXEC::sender auto multiple = STDEXEC::just() | exec::fork_join(then, then);
STATIC_REQUIRE(exec::sender_for<decltype(multiple), exec::fork_join_t>);
STATIC_REQUIRE(noexcept(exec::fork_join(STDEXEC::just(), then, then)));
}

struct ForwardingThen
{
template <typename Value>
Expand All @@ -69,6 +93,35 @@ namespace
};
#endif

//! Count how many copies and moves are performed.
struct counter
{
inline static std::atomic<unsigned int> copy_constructions{0};
inline static std::atomic<unsigned int> copy_assignments{0};
inline static std::atomic<unsigned int> move_constructions{0};
inline static std::atomic<unsigned int> move_assignments{0};

counter() = default;
counter(counter const &) noexcept
{
copy_constructions.fetch_add(1, std::memory_order_relaxed);
}
counter &operator=(counter const &) noexcept
{
copy_assignments.fetch_add(1, std::memory_order_relaxed);
return *this;
}
counter(counter &&) noexcept
{
move_constructions.fetch_add(1, std::memory_order_relaxed);
}
counter &operator=(counter &&) noexcept
{
move_assignments.fetch_add(1, std::memory_order_relaxed);
return *this;
}
};

template <char ID>
struct identifiable_domain : public STDEXEC::default_domain
{};
Expand Down Expand Up @@ -140,18 +193,43 @@ namespace
#if !STDEXEC_NO_STDCPP_EXCEPTIONS()
TEST_CASE("fork_join reports failures while caching results", "[adaptors][fork_join]")
{
std::atomic<int> witness{0};

auto sndr = exec::fork_join(exec::just_from(
[](auto sink) noexcept
{
static throwing_copy value;
return sink(value);
}),
then([](throwing_copy const &) noexcept {}));
then([&witness](throwing_copy const &) noexcept { ++witness; }),
then([&witness](throwing_copy const &) noexcept { ++witness; }));

CHECK_THROWS_AS(sync_wait(std::move(sndr)), int);

CHECK(witness == 0);
}
#endif

TEST_CASE("fork_join caches without copying when results are movable and replays by reference to "
"children",
"[adaptors][fork_join]")
{
std::atomic<int> witness{0};

auto sndr = exec::fork_join(exec::just_from([](auto sink) noexcept { return sink(counter{}); }),
STDEXEC::then([&witness](counter const &) noexcept { ++witness; }),
STDEXEC::then([&witness](counter const &) noexcept { ++witness; }));

STDEXEC::sync_wait(std::move(sndr));

CHECK(counter::copy_constructions == 0);
CHECK(counter::copy_assignments == 0);
CHECK(counter::move_constructions == 1);
CHECK(counter::move_assignments == 0);

CHECK(witness == 2);
}

TEST_CASE("fork_join can be nested", "[adaptors][fork_join]")
{
std::atomic<int> witness = 0;
Expand Down
11 changes: 11 additions & 0 deletions test/stdexec/algos/adaptors/test_then.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -179,6 +179,17 @@ namespace
| ex::then([] { return std::string{"hello"}; }));
}

TEST_CASE("then noexceptness", "[adaptors][then]")
{
auto func = [](){};
Comment thread
ericniebler marked this conversation as resolved.

STATIC_REQUIRE(noexcept(ex::then(func)));

STATIC_REQUIRE(noexcept(STDEXEC::just() | STDEXEC::then(func)));

STATIC_REQUIRE(noexcept(STDEXEC::then(STDEXEC::just(), func)));
}

TEST_CASE("then keeps error_types from input sender", "[adaptors][then]")
{
inline_scheduler sched1{};
Expand Down
Loading