Skip to content

Commit e172e8f

Browse files
author
Shivansh Singh
committed
Implement P2300 bulk adapter for HPX executors
- Add executor_scheduler_bulk.hpp with bulk sender/receiver for executor-based schedulers - Add executor_algorithm_bulk.cpp unit test covering sequenced and parallel bulk execution - Register new header and test in CMakeLists Signed-off-by: Shivansh Singh <your_github_email@example.com>
1 parent 2bb509d commit e172e8f

4 files changed

Lines changed: 269 additions & 0 deletions

File tree

libs/core/executors/CMakeLists.txt

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,7 @@ set(executors_headers
2727
hpx/executors/execution_policy_scheduling_property.hpp
2828
hpx/executors/execution_policy.hpp
2929
hpx/executors/executor_scheduler.hpp
30+
hpx/executors/executor_scheduler_bulk.hpp
3031
hpx/executors/explicit_scheduler_executor.hpp
3132
hpx/executors/fork_join_executor.hpp
3233
hpx/executors/limiting_executor.hpp
Lines changed: 173 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,173 @@
1+
// Copyright (c) 2026 The STE||AR-Group
2+
//
3+
// SPDX-License-Identifier: BSL-1.0
4+
// Distributed under the Boost Software License, Version 1.0. (See accompanying
5+
// file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
6+
7+
#pragma once
8+
9+
#include <hpx/config.hpp>
10+
11+
#include <hpx/execution_base/stdexec_forward.hpp>
12+
13+
#include <hpx/execution/algorithms/bulk.hpp>
14+
#include <hpx/execution_base/completion_scheduler.hpp>
15+
#include <hpx/execution_base/completion_signatures.hpp>
16+
#include <hpx/execution_base/execution.hpp>
17+
#include <hpx/execution_base/receiver.hpp>
18+
#include <hpx/execution_base/sender.hpp>
19+
#include <hpx/executors/executor_scheduler.hpp>
20+
#include <hpx/modules/errors.hpp>
21+
22+
#include <exception>
23+
#include <type_traits>
24+
#include <utility>
25+
26+
namespace hpx::execution::experimental {
27+
28+
namespace detail {
29+
template <typename Executor, typename Receiver, typename Shape,
30+
typename F>
31+
struct executor_bulk_receiver
32+
{
33+
using receiver_concept = hpx::execution::experimental::receiver_t;
34+
35+
HPX_NO_UNIQUE_ADDRESS std::decay_t<Executor> exec_;
36+
HPX_NO_UNIQUE_ADDRESS std::decay_t<Receiver> receiver_;
37+
HPX_NO_UNIQUE_ADDRESS std::decay_t<Shape> shape_;
38+
HPX_NO_UNIQUE_ADDRESS std::decay_t<F> f_;
39+
40+
template <typename Error>
41+
void set_error(Error&& error) noexcept
42+
{
43+
hpx::execution::experimental::set_error(
44+
HPX_MOVE(receiver_), HPX_FORWARD(Error, error));
45+
}
46+
47+
void set_stopped() noexcept
48+
{
49+
hpx::execution::experimental::set_stopped(HPX_MOVE(receiver_));
50+
}
51+
52+
template <typename... Ts>
53+
void set_value(Ts&&... ts) noexcept
54+
{
55+
hpx::detail::try_catch_exception_ptr(
56+
[&]() {
57+
hpx::parallel::execution::bulk_sync_execute(
58+
exec_, f_, shape_, ts...);
59+
60+
hpx::execution::experimental::set_value(
61+
HPX_MOVE(receiver_), HPX_FORWARD(Ts, ts)...);
62+
},
63+
[&](std::exception_ptr ep) {
64+
hpx::execution::experimental::set_error(
65+
HPX_MOVE(receiver_), HPX_MOVE(ep));
66+
});
67+
}
68+
};
69+
70+
template <typename Executor, typename Sender, typename Shape,
71+
typename F>
72+
struct executor_bulk_sender
73+
{
74+
HPX_NO_UNIQUE_ADDRESS std::decay_t<Executor> exec_;
75+
HPX_NO_UNIQUE_ADDRESS std::decay_t<Sender> sender_;
76+
HPX_NO_UNIQUE_ADDRESS std::decay_t<Shape> shape_;
77+
HPX_NO_UNIQUE_ADDRESS std::decay_t<F> f_;
78+
79+
using sender_concept = hpx::execution::experimental::sender_t;
80+
81+
template <typename Self, typename Env>
82+
static consteval auto get_completion_signatures() noexcept
83+
-> decltype(hpx::execution::experimental::
84+
transform_completion_signatures(
85+
hpx::execution::experimental::
86+
completion_signatures_of_t<Sender, Env>{},
87+
hpx::execution::experimental::keep_completion<
88+
hpx::execution::experimental::set_value_t>{},
89+
hpx::execution::experimental::keep_completion<
90+
hpx::execution::experimental::set_error_t>{},
91+
hpx::execution::experimental::keep_completion<
92+
hpx::execution::experimental::set_stopped_t>{},
93+
hpx::execution::experimental::completion_signatures<
94+
hpx::execution::experimental::set_error_t(
95+
std::exception_ptr)>{}))
96+
{
97+
return {};
98+
}
99+
100+
struct env
101+
{
102+
std::decay_t<Sender> const& pred_snd;
103+
std::decay_t<Executor> const& exec;
104+
105+
template <typename CPO>
106+
requires(
107+
meta::value<meta::one_of<CPO,
108+
hpx::execution::experimental::set_error_t,
109+
hpx::execution::experimental::set_stopped_t>> &&
110+
hpx::execution::experimental::detail::
111+
has_completion_scheduler_v<CPO,
112+
std::decay_t<Sender>>)
113+
constexpr auto query(
114+
hpx::execution::experimental::get_completion_scheduler_t<
115+
CPO>
116+
tag) const noexcept
117+
{
118+
return tag(hpx::execution::experimental::get_env(pred_snd));
119+
}
120+
121+
constexpr auto
122+
query(hpx::execution::experimental::get_completion_scheduler_t<
123+
hpx::execution::experimental::set_value_t>) const noexcept
124+
{
125+
return hpx::execution::experimental::executor_scheduler<
126+
Executor>{exec};
127+
}
128+
};
129+
130+
constexpr auto get_env() const noexcept
131+
{
132+
return env{sender_, exec_};
133+
}
134+
135+
template <typename Receiver>
136+
auto connect(Receiver&& receiver) &&
137+
{
138+
return hpx::execution::experimental::connect(HPX_MOVE(sender_),
139+
executor_bulk_receiver<Executor, std::decay_t<Receiver>,
140+
Shape, F>{HPX_MOVE(exec_),
141+
HPX_FORWARD(Receiver, receiver), HPX_MOVE(shape_),
142+
HPX_MOVE(f_)});
143+
}
144+
145+
template <typename Receiver>
146+
auto connect(Receiver&& receiver) &
147+
{
148+
return hpx::execution::experimental::connect(sender_,
149+
executor_bulk_receiver<Executor, std::decay_t<Receiver>,
150+
Shape, F>{
151+
exec_, HPX_FORWARD(Receiver, receiver), shape_, f_});
152+
}
153+
154+
template <typename Receiver>
155+
auto connect(Receiver&& receiver) const&
156+
{
157+
return hpx::execution::experimental::connect(sender_,
158+
executor_bulk_receiver<Executor, std::decay_t<Receiver>,
159+
Shape, F>{
160+
exec_, HPX_FORWARD(Receiver, receiver), shape_, f_});
161+
}
162+
};
163+
} // namespace detail
164+
165+
template <typename Executor, typename Sender, typename Shape, typename F>
166+
auto tag_invoke(bulk_t, executor_scheduler<Executor> const& sched,
167+
Sender&& sender, Shape const& shape, F&& f)
168+
{
169+
return detail::executor_bulk_sender<Executor, std::decay_t<Sender>,
170+
Shape, std::decay_t<F>>{
171+
sched.exec_, HPX_FORWARD(Sender, sender), shape, HPX_FORWARD(F, f)};
172+
}
173+
} // namespace hpx::execution::experimental

libs/core/executors/tests/unit/CMakeLists.txt

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@ set(tests
1010
created_executor
1111
execution_policy_mappings
1212
executor_scheduler
13+
executor_algorithm_bulk
1314
explicit_scheduler_executor
1415
fork_join_executor
1516
fork_join_executor_from
Lines changed: 94 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,94 @@
1+
// Copyright (c) 2026 The STE||AR-Group
2+
//
3+
// SPDX-License-Identifier: BSL-1.0
4+
// Distributed under the Boost Software License, Version 1.0. (See accompanying
5+
// file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
6+
7+
#include <hpx/execution.hpp>
8+
#include <hpx/init.hpp>
9+
#include <hpx/modules/testing.hpp>
10+
11+
#include <hpx/modules/executors.hpp>
12+
13+
#include <atomic>
14+
#include <cstdint>
15+
#include <iostream>
16+
#include <string>
17+
#include <vector>
18+
19+
namespace ex = hpx::execution::experimental;
20+
21+
void test_sequential_bulk()
22+
{
23+
hpx::execution::sequenced_executor exec;
24+
std::atomic<int> call_count{0};
25+
26+
// Obtain a P2300 scheduler from the executor via query(get_scheduler_t{})
27+
auto sched = exec.query(ex::get_scheduler_t{});
28+
29+
auto snd = ex::schedule(sched) | ex::bulk(10, [&](int i) {
30+
(void) i;
31+
++call_count;
32+
// Should run sequentially, so it's safe
33+
});
34+
35+
hpx::this_thread::experimental::sync_wait(snd);
36+
37+
HPX_TEST_EQ(call_count.load(), 10);
38+
}
39+
40+
void test_parallel_bulk()
41+
{
42+
hpx::execution::parallel_executor exec;
43+
std::atomic<int> call_count{0};
44+
45+
auto sched = exec.query(ex::get_scheduler_t{});
46+
47+
auto snd = ex::schedule(sched) | ex::bulk(1000, [&](int i) {
48+
(void) i;
49+
++call_count;
50+
});
51+
52+
hpx::this_thread::experimental::sync_wait(snd);
53+
54+
HPX_TEST_EQ(call_count.load(), 1000);
55+
}
56+
57+
void test_parallel_bulk_with_value()
58+
{
59+
hpx::execution::parallel_executor exec;
60+
std::atomic<int> call_count{0};
61+
62+
auto sched = exec.query(ex::get_scheduler_t{});
63+
64+
auto snd = ex::schedule(sched) | ex::then([]() { return 42; }) |
65+
ex::bulk(100, [&](int i, int val) {
66+
(void) i;
67+
HPX_TEST_EQ(val, 42);
68+
++call_count;
69+
});
70+
71+
hpx::this_thread::experimental::sync_wait(snd);
72+
73+
HPX_TEST_EQ(call_count.load(), 100);
74+
}
75+
76+
int hpx_main()
77+
{
78+
test_sequential_bulk();
79+
test_parallel_bulk();
80+
test_parallel_bulk_with_value();
81+
return hpx::local::finalize();
82+
}
83+
84+
int main(int argc, char* argv[])
85+
{
86+
std::vector<std::string> const cfg = {"hpx.os_threads=all"};
87+
hpx::local::init_params init_args;
88+
init_args.cfg = cfg;
89+
90+
HPX_TEST_EQ_MSG(hpx::local::init(hpx_main, argc, argv, init_args), 0,
91+
"HPX main exited with non-zero status");
92+
93+
return hpx::util::report_errors();
94+
}

0 commit comments

Comments
 (0)