Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
26 commits
Select commit Hold shift + click to select a range
61e498d
feat(execution): implement HPX future-sender bridge
shivansh023023 Jun 21, 2026
61d7fa9
fix(execution): add missing file for future-sender bridge
shivansh023023 Jun 23, 2026
55d4144
fix(execution): resolve C++20 modules CI and hpxinspect failures
shivansh023023 Jul 2, 2026
b5a7a8d
test(execution): fix thread_pool_scheduler tests and deadlock in futu…
shivansh023023 Jul 4, 2026
c5a1b25
chore: trigger CI
shivansh023023 Jul 4, 2026
00a7bb2
fix(execution): apply CodeRabbit fixes and Clang-22 CI workaround
shivansh023023 Jul 5, 2026
e840440
fix(ci): isolate Clang module compilation ICE to future_sender_test
shivansh023023 Jul 6, 2026
ed9b453
fix(futures): resolve implicit instantiation error in sender_future
shivansh023023 Jul 6, 2026
9873c27
refactor(futures): cleanup sender_future templates, optimize allocati…
shivansh023023 Jul 12, 2026
9af78ed
fix(futures): resolve as_sender CPO redefinition by migrating to tag_…
shivansh023023 Jul 12, 2026
8b27a8f
fix(futures): correct as_sender include order to fix CI compilation
shivansh023023 Jul 12, 2026
d18a802
fix(futures): break circular dependency on execution module by forwar…
shivansh023023 Jul 12, 2026
7381b26
refactor(execution): migrate future/sender bridges to execution modul…
shivansh023023 Jul 12, 2026
1989260
fix(execution): resolve CodeRabbit review for future_sender docs and …
shivansh023023 Jul 14, 2026
779f0c7
execution: route future sender bridge through CPOs
guptapratykshh Jul 15, 2026
ce839d0
Merge pull request #1 from guptapratykshh/feat/as-future-cpo
shivansh023023 Jul 19, 2026
1d2c389
execution: optimize future sender scheduling bridge
guptapratykshh Jul 20, 2026
bd236ca
execution: add future sender bridge benchmark
guptapratykshh Jul 20, 2026
0af562d
execution: keep future sender scheduler bridge lightweight
guptapratykshh Jul 20, 2026
174549d
execution: add missing future sender test include
guptapratykshh Jul 20, 2026
a1c8d28
Merge pull request #2 from guptapratykshh/feat/as-sender-continues-on…
shivansh023023 Jul 20, 2026
03683cd
fix(execution): address CodeRabbit review feedback on sender bridge
shivansh023023 Jul 20, 2026
ea687bc
execution: fix future sender benchmark comment
guptapratykshh Jul 21, 2026
606bc42
Merge pull request #4 from guptapratykshh/fix/7256-future-sender-comment
shivansh023023 Jul 21, 2026
72fdb78
execution: avoid deprecated volatile compound assignment
guptapratykshh Jul 28, 2026
6051377
Merge pull request #5 from guptapratykshh/fix/7256-volatile-benchmark
shivansh023023 Jul 30, 2026
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
2 changes: 2 additions & 0 deletions libs/core/execution/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -15,9 +15,11 @@ set(execution_headers
hpx/execution/algorithms/detail/predicates.hpp
hpx/execution/algorithms/detail/single_result.hpp
hpx/execution/algorithms/detail/sync_wait_domain.hpp
hpx/execution/algorithms/future_sender.hpp
hpx/execution/algorithms/keep_future.hpp
hpx/execution/algorithms/make_future.hpp
hpx/execution/algorithms/run_loop.hpp
hpx/execution/algorithms/sender_future.hpp
hpx/execution/algorithms/sync_wait.hpp
hpx/execution/algorithms/when_all.hpp
hpx/execution/algorithms/when_all_vector.hpp
Expand Down
317 changes: 6 additions & 311 deletions libs/core/execution/include/hpx/execution/algorithms/as_sender.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@

#include <hpx/assert.hpp>
#include <hpx/execution/algorithms/detail/partial_algorithm.hpp>
#include <hpx/execution/algorithms/future_sender.hpp>
#include <hpx/modules/concepts.hpp>
#include <hpx/modules/errors.hpp>
#include <hpx/modules/execution_base.hpp>
Expand All @@ -18,312 +19,6 @@
#include <type_traits>
#include <utility>

namespace hpx::execution::experimental { namespace detail {

///////////////////////////////////////////////////////////////////////////
// Operation state for sender compatibility
HPX_CXX_CORE_EXPORT template <typename Receiver, typename Future>
class as_sender_operation_state
{
private:
using receiver_type = std::decay_t<Receiver>;
using future_type = std::decay_t<Future>;
using result_type = typename future_type::result_type;

public:
template <typename Receiver_>
as_sender_operation_state(Receiver_&& r, future_type f)
: receiver_(HPX_FORWARD(Receiver_, r))
, future_(HPX_MOVE(f))
{
}

as_sender_operation_state(as_sender_operation_state&&) = delete;
as_sender_operation_state& operator=(
as_sender_operation_state&&) = delete;
as_sender_operation_state(as_sender_operation_state const&) = delete;
as_sender_operation_state& operator=(
as_sender_operation_state const&) = delete;

void start() & noexcept
{
start_helper();
}

private:
void start_helper() & noexcept
{
hpx::detail::try_catch_exception_ptr(
[&]() {
auto state = traits::detail::get_shared_state(future_);

if (!state)
{
HPX_THROW_EXCEPTION(hpx::error::no_state,
"as_sender_operation_state::start",
"the future has no valid shared state");
}

auto on_completed = [this]() mutable {
if (future_.has_value())
{
if constexpr (std::is_void_v<result_type>)
{
hpx::execution::experimental::set_value(
HPX_MOVE(receiver_));
}
else
{
hpx::execution::experimental::set_value(
HPX_MOVE(receiver_), future_.get());
}
}
else if (future_.has_exception())
{
hpx::execution::experimental::set_error(
HPX_MOVE(receiver_),
future_.get_exception_ptr());
}
};

if (!state->is_ready(std::memory_order_relaxed))
{
state->execute_deferred();

// execute_deferred might have made the future ready
if (!state->is_ready(std::memory_order_relaxed))
{
// The operation state has to be kept alive until
// set_value is called, which means that we don't
// need to move receiver and future into the
// on_completed callback.
state->set_on_completed(HPX_MOVE(on_completed));
}
else
{
on_completed();
}
}
else
{
on_completed();
}
},
[&](std::exception_ptr ep) {
hpx::execution::experimental::set_error(
HPX_MOVE(receiver_), HPX_MOVE(ep));
});
}

HPX_NO_UNIQUE_ADDRESS std::decay_t<Receiver> receiver_;
future_type future_;
};

HPX_CXX_CORE_EXPORT template <typename Future>
struct as_sender_sender_base
{
using result_type = typename std::decay_t<Future>::result_type;

std::decay_t<Future> future_;

template <bool IsVoid, typename _result_type>
struct set_value_void_checked
{
using type = hpx::execution::experimental::set_value_t(
_result_type);
};

template <typename _result_type>
struct set_value_void_checked<true, _result_type>
{
using type = hpx::execution::experimental::set_value_t();
};

using completion_signatures =
hpx::execution::experimental::completion_signatures<
typename set_value_void_checked<std::is_void_v<result_type>,
result_type>::type,
hpx::execution::experimental::set_error_t(std::exception_ptr)>;
};

HPX_CXX_CORE_EXPORT template <typename Future>
struct as_sender_sender;

template <typename T>
struct as_sender_sender<hpx::future<T>>
: public as_sender_sender_base<hpx::future<T>>
{
using sender_concept = hpx::execution::experimental::sender_t;
using future_type = hpx::future<T>;
using base_type = as_sender_sender_base<hpx::future<T>>;
using base_type::future_;

template <typename Future>
requires(!std::is_same_v<std::decay_t<Future>, as_sender_sender>)
explicit as_sender_sender(Future&& future)
: base_type{HPX_FORWARD(Future, future)}
{
}

as_sender_sender(as_sender_sender&&) = default;
as_sender_sender& operator=(as_sender_sender&&) = default;
as_sender_sender(as_sender_sender const&) = delete;
as_sender_sender& operator=(as_sender_sender const&) = delete;

template <typename Self, typename... Env>
static consteval auto get_completion_signatures() noexcept ->
typename base_type::completion_signatures
{
return {};
}

template <typename Receiver>
auto connect(Receiver&& receiver) &&
{
return as_sender_operation_state<Receiver, future_type>{
HPX_FORWARD(Receiver, receiver), HPX_MOVE(future_)};
}
};
}} // namespace hpx::execution::experimental::detail

namespace hpx::execution::experimental { namespace detail {
template <typename T>
struct as_sender_sender<hpx::shared_future<T>>
: as_sender_sender_base<hpx::shared_future<T>>
{
using sender_concept = hpx::execution::experimental::sender_t;
using future_type = hpx::shared_future<T>;
using base_type = as_sender_sender_base<hpx::shared_future<T>>;
using base_type::future_;

template <typename Future>
requires(!std::is_same_v<std::decay_t<Future>, as_sender_sender>)
explicit as_sender_sender(Future&& future)
: base_type{HPX_FORWARD(Future, future)}
{
}

as_sender_sender(as_sender_sender&&) = default;
as_sender_sender& operator=(as_sender_sender&&) = default;
as_sender_sender(as_sender_sender const&) = default;
as_sender_sender& operator=(as_sender_sender const&) = default;

template <typename Self, typename... Env>
static consteval auto get_completion_signatures() noexcept ->
typename base_type::completion_signatures
{
return {};
}

template <typename Receiver>
auto connect(Receiver&& receiver) &&
{
return as_sender_operation_state<Receiver, future_type>{
HPX_FORWARD(Receiver, receiver), HPX_MOVE(future_)};
}

template <typename Receiver>
auto connect(Receiver&& receiver) &
{
return as_sender_operation_state<Receiver, future_type>{
HPX_FORWARD(Receiver, receiver), future_};
}
};
}} // namespace hpx::execution::experimental::detail

namespace hpx::execution::experimental { namespace detail {

///////////////////////////////////////////////////////////////////////
// Scheduler-aware sender wrapper.
//
// Exposes the stored scheduler through its environment so that
// downstream sender algorithms (bulk, sync_wait, etc.) can query
// get_completion_scheduler<set_value_t> and obtain the scheduler
// that originated the work.
HPX_CXX_CORE_EXPORT template <typename Future, typename Scheduler>
requires(hpx::traits::is_future_v<std::decay_t<Future>>)
struct as_sender_sender_with_scheduler
: public as_sender_sender_base<std::decay_t<Future>>
{
using sender_concept = hpx::execution::experimental::sender_t;
using future_type = std::decay_t<Future>;
using scheduler_type = std::decay_t<Scheduler>;
using base_type = as_sender_sender_base<std::decay_t<Future>>;
using base_type::future_;

HPX_NO_UNIQUE_ADDRESS scheduler_type scheduler_;

// Environment that answers get_completion_scheduler queries.
struct env
{
scheduler_type sched;

auto query(
hpx::execution::experimental::get_domain_t) const noexcept
{
return hpx::execution::experimental::get_domain(sched);
}

template <typename CPO>
requires(std::is_same_v<CPO,
hpx::execution::experimental::set_value_t> ||
std::is_same_v<CPO,
hpx::execution::experimental::set_stopped_t>)
auto query(
hpx::execution::experimental::get_completion_scheduler_t<CPO>)
const noexcept
{
return sched;
}
};

template <typename Future_, typename Scheduler_>
requires(!std::is_same_v<std::decay_t<Future_>,
as_sender_sender_with_scheduler>)
explicit as_sender_sender_with_scheduler(
Future_&& future, Scheduler_&& scheduler)
: base_type{HPX_FORWARD(Future_, future)}
, scheduler_(HPX_FORWARD(Scheduler_, scheduler))
{
}

as_sender_sender_with_scheduler(
as_sender_sender_with_scheduler&&) = default;
as_sender_sender_with_scheduler& operator=(
as_sender_sender_with_scheduler&&) = default;
as_sender_sender_with_scheduler(
as_sender_sender_with_scheduler const&) = default;
as_sender_sender_with_scheduler& operator=(
as_sender_sender_with_scheduler const&) = default;

template <typename Self, typename... Env>
static consteval auto get_completion_signatures() noexcept ->
typename base_type::completion_signatures
{
return {};
}

template <typename Receiver>
auto connect(Receiver&& receiver) &&
{
return as_sender_operation_state<Receiver, future_type>{
HPX_FORWARD(Receiver, receiver), HPX_MOVE(future_)};
}

template <typename Receiver>
auto connect(Receiver&& receiver) &
{
return as_sender_operation_state<Receiver, future_type>{
HPX_FORWARD(Receiver, receiver), future_};
}

constexpr auto get_env() const noexcept
{
return env{scheduler_};
}
};
}} // namespace hpx::execution::experimental::detail

namespace hpx::execution::experimental {
// The as_sender CPO can be used to adapt any HPX future as a sender. The
// value provided by the future will be used to call set_value on the
Expand All @@ -343,8 +38,8 @@ namespace hpx::execution::experimental {
// clang-format on
constexpr HPX_FORCEINLINE auto operator()(Future&& future) const
{
return detail::as_sender_sender<std::decay_t<Future>>(
HPX_FORWARD(Future, future));
return detail::future_sender<std::decay_t<Future>>{
HPX_FORWARD(Future, future)};
}

// Scheduler-aware overload: wraps the future into a sender whose
Expand All @@ -356,9 +51,9 @@ namespace hpx::execution::experimental {
constexpr HPX_FORCEINLINE auto operator()(
Future&& future, Scheduler&& scheduler) const
{
return detail::as_sender_sender_with_scheduler<std::decay_t<Future>,
std::decay_t<Scheduler>>(
HPX_FORWARD(Future, future), HPX_FORWARD(Scheduler, scheduler));
return detail::future_sender_with_scheduler<std::decay_t<Future>,
std::decay_t<Scheduler>>{
HPX_FORWARD(Future, future), HPX_FORWARD(Scheduler, scheduler)};
}

constexpr HPX_FORCEINLINE auto operator()() const
Expand Down
Loading
Loading