parallel_group.hpp 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434
  1. //
  2. // experimental/impl/parallel_group.hpp
  3. // ~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
  4. //
  5. // Copyright (c) 2003-2022 Christopher M. Kohlhoff (chris at kohlhoff dot com)
  6. //
  7. // Distributed under the Boost Software License, Version 1.0. (See accompanying
  8. // file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
  9. //
  10. #ifndef ASIO_IMPL_EXPERIMENTAL_PARALLEL_GROUP_HPP
  11. #define ASIO_IMPL_EXPERIMENTAL_PARALLEL_GROUP_HPP
  12. #if defined(_MSC_VER) && (_MSC_VER >= 1200)
  13. # pragma once
  14. #endif // defined(_MSC_VER) && (_MSC_VER >= 1200)
  15. #include "asio/detail/config.hpp"
  16. #include <atomic>
  17. #include <memory>
  18. #include <new>
  19. #include <tuple>
  20. #include "asio/associated_cancellation_slot.hpp"
  21. #include "asio/detail/recycling_allocator.hpp"
  22. #include "asio/detail/type_traits.hpp"
  23. #include "asio/dispatch.hpp"
  24. #include "asio/detail/push_options.hpp"
  25. namespace asio {
  26. namespace experimental {
  27. namespace detail {
  28. // Stores the result from an individual asynchronous operation.
  29. template <typename T, typename = void>
  30. struct parallel_group_op_result
  31. {
  32. public:
  33. parallel_group_op_result()
  34. : has_value_(false)
  35. {
  36. }
  37. parallel_group_op_result(parallel_group_op_result&& other)
  38. : has_value_(other.has_value_)
  39. {
  40. if (has_value_)
  41. new (&u_.value_) T(std::move(other.get()));
  42. }
  43. ~parallel_group_op_result()
  44. {
  45. if (has_value_)
  46. u_.value_.~T();
  47. }
  48. T& get() noexcept
  49. {
  50. return u_.value_;
  51. }
  52. template <typename... Args>
  53. void emplace(Args&&... args)
  54. {
  55. new (&u_.value_) T(std::forward<Args>(args)...);
  56. has_value_ = true;
  57. }
  58. private:
  59. union u
  60. {
  61. u() {}
  62. ~u() {}
  63. char c_;
  64. T value_;
  65. } u_;
  66. bool has_value_;
  67. };
  68. // Proxy completion handler for the group of parallel operatations. Unpacks and
  69. // concatenates the individual operations' results, and invokes the user's
  70. // completion handler.
  71. template <typename Handler, typename... Ops>
  72. struct parallel_group_completion_handler
  73. {
  74. typedef typename decay<
  75. typename prefer_result<
  76. typename associated_executor<Handler>::type,
  77. execution::outstanding_work_t::tracked_t
  78. >::type
  79. >::type executor_type;
  80. parallel_group_completion_handler(Handler&& h)
  81. : handler_(std::move(h)),
  82. executor_(
  83. asio::prefer(
  84. asio::get_associated_executor(handler_),
  85. execution::outstanding_work.tracked))
  86. {
  87. }
  88. executor_type get_executor() const noexcept
  89. {
  90. return executor_;
  91. }
  92. void operator()()
  93. {
  94. this->invoke(std::make_index_sequence<sizeof...(Ops)>());
  95. }
  96. template <std::size_t... I>
  97. void invoke(std::index_sequence<I...>)
  98. {
  99. this->invoke(std::tuple_cat(std::move(std::get<I>(args_).get())...));
  100. }
  101. template <typename... Args>
  102. void invoke(std::tuple<Args...>&& args)
  103. {
  104. this->invoke(std::move(args), std::make_index_sequence<sizeof...(Args)>());
  105. }
  106. template <typename... Args, std::size_t... I>
  107. void invoke(std::tuple<Args...>&& args, std::index_sequence<I...>)
  108. {
  109. std::move(handler_)(completion_order_, std::move(std::get<I>(args))...);
  110. }
  111. Handler handler_;
  112. executor_type executor_;
  113. std::array<std::size_t, sizeof...(Ops)> completion_order_{};
  114. std::tuple<
  115. parallel_group_op_result<
  116. typename parallel_op_signature_as_tuple<
  117. typename parallel_op_signature<Ops>::type
  118. >::type
  119. >...
  120. > args_{};
  121. };
  122. // Shared state for the parallel group.
  123. template <typename Condition, typename Handler, typename... Ops>
  124. struct parallel_group_state
  125. {
  126. parallel_group_state(Condition&& c, Handler&& h)
  127. : cancellation_condition_(std::move(c)),
  128. handler_(std::move(h))
  129. {
  130. }
  131. // The number of operations that have completed so far. Used to determine the
  132. // order of completion.
  133. std::atomic<unsigned int> completed_{0};
  134. // The non-none cancellation type that resulted from a cancellation condition.
  135. // Stored here for use by the group's initiating function.
  136. std::atomic<cancellation_type_t> cancel_type_{cancellation_type::none};
  137. // The number of cancellations that have been requested, either on completion
  138. // of the operations within the group, or via the cancellation slot for the
  139. // group operation. Initially set to the number of operations to prevent
  140. // cancellation signals from being emitted until after all of the group's
  141. // operations' initiating functions have completed.
  142. std::atomic<unsigned int> cancellations_requested_{sizeof...(Ops)};
  143. // The number of operations that are yet to complete. Used to determine when
  144. // it is safe to invoke the user's completion handler.
  145. std::atomic<unsigned int> outstanding_{sizeof...(Ops)};
  146. // The cancellation signals for each operation in the group.
  147. asio::cancellation_signal cancellation_signals_[sizeof...(Ops)];
  148. // The cancellation condition is used to determine whether the results from an
  149. // individual operation warrant a cancellation request for the whole group.
  150. Condition cancellation_condition_;
  151. // The proxy handler to be invoked once all operations in the group complete.
  152. parallel_group_completion_handler<Handler, Ops...> handler_;
  153. };
  154. // Handler for an individual operation within the parallel group.
  155. template <std::size_t I, typename Condition, typename Handler, typename... Ops>
  156. struct parallel_group_op_handler
  157. {
  158. typedef asio::cancellation_slot cancellation_slot_type;
  159. parallel_group_op_handler(
  160. std::shared_ptr<parallel_group_state<Condition, Handler, Ops...> > state)
  161. : state_(std::move(state))
  162. {
  163. }
  164. cancellation_slot_type get_cancellation_slot() const noexcept
  165. {
  166. return state_->cancellation_signals_[I].slot();
  167. }
  168. template <typename... Args>
  169. void operator()(Args... args)
  170. {
  171. // Capture this operation into the completion order.
  172. state_->handler_.completion_order_[state_->completed_++] = I;
  173. // Determine whether the results of this operation require cancellation of
  174. // the whole group.
  175. cancellation_type_t cancel_type = state_->cancellation_condition_(args...);
  176. // Capture the result of the operation into the proxy completion handler.
  177. std::get<I>(state_->handler_.args_).emplace(std::move(args)...);
  178. if (cancel_type != cancellation_type::none)
  179. {
  180. // Save the type for potential use by the group's initiating function.
  181. state_->cancel_type_ = cancel_type;
  182. // If we are the first operation to request cancellation, emit a signal
  183. // for each operation in the group.
  184. if (state_->cancellations_requested_++ == 0)
  185. for (std::size_t i = 0; i < sizeof...(Ops); ++i)
  186. if (i != I)
  187. state_->cancellation_signals_[i].emit(cancel_type);
  188. }
  189. // If this is the last outstanding operation, invoke the user's handler.
  190. if (--state_->outstanding_ == 0)
  191. asio::dispatch(std::move(state_->handler_));
  192. }
  193. std::shared_ptr<parallel_group_state<Condition, Handler, Ops...> > state_;
  194. };
  195. // Handler for an individual operation within the parallel group that has an
  196. // explicitly specified executor.
  197. template <typename Executor, std::size_t I,
  198. typename Condition, typename Handler, typename... Ops>
  199. struct parallel_group_op_handler_with_executor :
  200. parallel_group_op_handler<I, Condition, Handler, Ops...>
  201. {
  202. typedef parallel_group_op_handler<I, Condition, Handler, Ops...> base_type;
  203. typedef asio::cancellation_slot cancellation_slot_type;
  204. typedef Executor executor_type;
  205. parallel_group_op_handler_with_executor(
  206. std::shared_ptr<parallel_group_state<Condition, Handler, Ops...> > state,
  207. executor_type ex)
  208. : parallel_group_op_handler<I, Condition, Handler, Ops...>(std::move(state))
  209. {
  210. cancel_proxy_ =
  211. &this->state_->cancellation_signals_[I].slot().template
  212. emplace<cancel_proxy>(this->state_, std::move(ex));
  213. }
  214. cancellation_slot_type get_cancellation_slot() const noexcept
  215. {
  216. return cancel_proxy_->signal_.slot();
  217. }
  218. executor_type get_executor() const noexcept
  219. {
  220. return cancel_proxy_->executor_;
  221. }
  222. // Proxy handler that forwards the emitted signal to the correct executor.
  223. struct cancel_proxy
  224. {
  225. cancel_proxy(
  226. std::shared_ptr<parallel_group_state<
  227. Condition, Handler, Ops...> > state,
  228. executor_type ex)
  229. : state_(std::move(state)),
  230. executor_(std::move(ex))
  231. {
  232. }
  233. void operator()(cancellation_type_t type)
  234. {
  235. if (auto state = state_.lock())
  236. {
  237. asio::cancellation_signal* sig = &signal_;
  238. asio::dispatch(executor_,
  239. [state, sig, type]{ sig->emit(type); });
  240. }
  241. }
  242. std::weak_ptr<parallel_group_state<Condition, Handler, Ops...> > state_;
  243. asio::cancellation_signal signal_;
  244. executor_type executor_;
  245. };
  246. cancel_proxy* cancel_proxy_;
  247. };
  248. // Helper to launch an operation using the correct executor, if any.
  249. template <std::size_t I, typename Op, typename = void>
  250. struct parallel_group_op_launcher
  251. {
  252. template <typename Condition, typename Handler, typename... Ops>
  253. static void launch(Op& op,
  254. const std::shared_ptr<parallel_group_state<
  255. Condition, Handler, Ops...> >& state)
  256. {
  257. typedef typename associated_executor<Op>::type ex_type;
  258. ex_type ex = asio::get_associated_executor(op);
  259. std::move(op)(
  260. parallel_group_op_handler_with_executor<ex_type, I,
  261. Condition, Handler, Ops...>(state, std::move(ex)));
  262. }
  263. };
  264. // Specialised launcher for operations that specify no executor.
  265. template <std::size_t I, typename Op>
  266. struct parallel_group_op_launcher<I, Op,
  267. typename enable_if<
  268. is_same<
  269. typename associated_executor<
  270. Op>::asio_associated_executor_is_unspecialised,
  271. void
  272. >::value
  273. >::type>
  274. {
  275. template <typename Condition, typename Handler, typename... Ops>
  276. static void launch(Op& op,
  277. const std::shared_ptr<parallel_group_state<
  278. Condition, Handler, Ops...> >& state)
  279. {
  280. std::move(op)(
  281. parallel_group_op_handler<I, Condition, Handler, Ops...>(state));
  282. }
  283. };
  284. template <typename Condition, typename Handler, typename... Ops>
  285. struct parallel_group_cancellation_handler
  286. {
  287. parallel_group_cancellation_handler(
  288. std::shared_ptr<parallel_group_state<Condition, Handler, Ops...> > state)
  289. : state_(std::move(state))
  290. {
  291. }
  292. void operator()(cancellation_type_t cancel_type)
  293. {
  294. // If we are the first place to request cancellation, i.e. no operation has
  295. // yet completed and requested cancellation, emit a signal for each
  296. // operation in the group.
  297. if (cancel_type != cancellation_type::none)
  298. if (auto state = state_.lock())
  299. if (state->cancellations_requested_++ == 0)
  300. for (std::size_t i = 0; i < sizeof...(Ops); ++i)
  301. state->cancellation_signals_[i].emit(cancel_type);
  302. }
  303. std::weak_ptr<parallel_group_state<Condition, Handler, Ops...> > state_;
  304. };
  305. template <typename Condition, typename Handler,
  306. typename... Ops, std::size_t... I>
  307. void parallel_group_launch(Condition cancellation_condition, Handler handler,
  308. std::tuple<Ops...>& ops, std::index_sequence<I...>)
  309. {
  310. // Get the user's completion handler's cancellation slot, so that we can allow
  311. // cancellation of the entire group.
  312. typename associated_cancellation_slot<Handler>::type slot
  313. = asio::get_associated_cancellation_slot(handler);
  314. // Create the shared state for the operation.
  315. typedef parallel_group_state<Condition, Handler, Ops...> state_type;
  316. std::shared_ptr<state_type> state = std::allocate_shared<state_type>(
  317. asio::detail::recycling_allocator<state_type,
  318. asio::detail::thread_info_base::parallel_group_tag>(),
  319. std::move(cancellation_condition), std::move(handler));
  320. // Initiate each individual operation in the group.
  321. int fold[] = { 0,
  322. ( parallel_group_op_launcher<I, Ops>::launch(std::get<I>(ops), state),
  323. 0 )...
  324. };
  325. (void)fold;
  326. // Check if any of the operations has already requested cancellation, and if
  327. // so, emit a signal for each operation in the group.
  328. if ((state->cancellations_requested_ -= sizeof...(Ops)) > 0)
  329. for (auto& signal : state->cancellation_signals_)
  330. signal.emit(state->cancel_type_);
  331. // Register a handler with the user's completion handler's cancellation slot.
  332. if (slot.is_connected())
  333. slot.template emplace<
  334. parallel_group_cancellation_handler<
  335. Condition, Handler, Ops...> >(state);
  336. }
  337. } // namespace detail
  338. } // namespace experimental
  339. template <typename R, typename... Args>
  340. class async_result<
  341. experimental::detail::parallel_op_signature_probe,
  342. R(Args...)>
  343. {
  344. public:
  345. typedef experimental::detail::parallel_op_signature_probe_result<
  346. void(Args...)> return_type;
  347. template <typename Initiation, typename... InitArgs>
  348. static return_type initiate(Initiation&&,
  349. experimental::detail::parallel_op_signature_probe, InitArgs&&...)
  350. {
  351. return return_type{};
  352. }
  353. };
  354. template <template <typename, typename> class Associator,
  355. typename Handler, typename... Ops, typename DefaultCandidate>
  356. struct associator<Associator,
  357. experimental::detail::parallel_group_completion_handler<Handler, Ops...>,
  358. DefaultCandidate>
  359. : Associator<Handler, DefaultCandidate>
  360. {
  361. static typename Associator<Handler, DefaultCandidate>::type get(
  362. const experimental::detail::parallel_group_completion_handler<
  363. Handler, Ops...>& h,
  364. const DefaultCandidate& c = DefaultCandidate()) ASIO_NOEXCEPT
  365. {
  366. return Associator<Handler, DefaultCandidate>::get(h.handler_, c);
  367. }
  368. };
  369. } // namespace asio
  370. #include "asio/detail/pop_options.hpp"
  371. #endif // ASIO_IMPL_EXPERIMENTAL_PARALLEL_GROUP_HPP