thread_pool.hpp 10 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355
  1. //
  2. // impl/thread_pool.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_THREAD_POOL_HPP
  11. #define ASIO_IMPL_THREAD_POOL_HPP
  12. #if defined(_MSC_VER) && (_MSC_VER >= 1200)
  13. # pragma once
  14. #endif // defined(_MSC_VER) && (_MSC_VER >= 1200)
  15. #include "asio/detail/blocking_executor_op.hpp"
  16. #include "asio/detail/bulk_executor_op.hpp"
  17. #include "asio/detail/executor_op.hpp"
  18. #include "asio/detail/fenced_block.hpp"
  19. #include "asio/detail/non_const_lvalue.hpp"
  20. #include "asio/detail/type_traits.hpp"
  21. #include "asio/execution_context.hpp"
  22. #include "asio/detail/push_options.hpp"
  23. namespace asio {
  24. inline thread_pool::executor_type
  25. thread_pool::get_executor() ASIO_NOEXCEPT
  26. {
  27. return executor_type(*this);
  28. }
  29. inline thread_pool::executor_type
  30. thread_pool::executor() ASIO_NOEXCEPT
  31. {
  32. return executor_type(*this);
  33. }
  34. inline thread_pool::scheduler_type
  35. thread_pool::scheduler() ASIO_NOEXCEPT
  36. {
  37. return scheduler_type(*this);
  38. }
  39. template <typename Allocator, unsigned int Bits>
  40. thread_pool::basic_executor_type<Allocator, Bits>&
  41. thread_pool::basic_executor_type<Allocator, Bits>::operator=(
  42. const basic_executor_type& other) ASIO_NOEXCEPT
  43. {
  44. if (this != &other)
  45. {
  46. thread_pool* old_thread_pool = pool_;
  47. pool_ = other.pool_;
  48. allocator_ = other.allocator_;
  49. bits_ = other.bits_;
  50. if (Bits & outstanding_work_tracked)
  51. {
  52. if (pool_)
  53. pool_->scheduler_.work_started();
  54. if (old_thread_pool)
  55. old_thread_pool->scheduler_.work_finished();
  56. }
  57. }
  58. return *this;
  59. }
  60. #if defined(ASIO_HAS_MOVE)
  61. template <typename Allocator, unsigned int Bits>
  62. thread_pool::basic_executor_type<Allocator, Bits>&
  63. thread_pool::basic_executor_type<Allocator, Bits>::operator=(
  64. basic_executor_type&& other) ASIO_NOEXCEPT
  65. {
  66. if (this != &other)
  67. {
  68. thread_pool* old_thread_pool = pool_;
  69. pool_ = other.pool_;
  70. allocator_ = std::move(other.allocator_);
  71. bits_ = other.bits_;
  72. if (Bits & outstanding_work_tracked)
  73. {
  74. other.pool_ = 0;
  75. if (old_thread_pool)
  76. old_thread_pool->scheduler_.work_finished();
  77. }
  78. }
  79. return *this;
  80. }
  81. #endif // defined(ASIO_HAS_MOVE)
  82. template <typename Allocator, unsigned int Bits>
  83. inline bool thread_pool::basic_executor_type<Allocator,
  84. Bits>::running_in_this_thread() const ASIO_NOEXCEPT
  85. {
  86. return pool_->scheduler_.can_dispatch();
  87. }
  88. template <typename Allocator, unsigned int Bits>
  89. template <typename Function>
  90. void thread_pool::basic_executor_type<Allocator,
  91. Bits>::do_execute(ASIO_MOVE_ARG(Function) f, false_type) const
  92. {
  93. typedef typename decay<Function>::type function_type;
  94. // Invoke immediately if the blocking.possibly property is enabled and we are
  95. // already inside the thread pool.
  96. if ((bits_ & blocking_never) == 0 && pool_->scheduler_.can_dispatch())
  97. {
  98. // Make a local, non-const copy of the function.
  99. function_type tmp(ASIO_MOVE_CAST(Function)(f));
  100. #if defined(ASIO_HAS_STD_EXCEPTION_PTR) \
  101. && !defined(ASIO_NO_EXCEPTIONS)
  102. try
  103. {
  104. #endif // defined(ASIO_HAS_STD_EXCEPTION_PTR)
  105. // && !defined(ASIO_NO_EXCEPTIONS)
  106. detail::fenced_block b(detail::fenced_block::full);
  107. asio_handler_invoke_helpers::invoke(tmp, tmp);
  108. return;
  109. #if defined(ASIO_HAS_STD_EXCEPTION_PTR) \
  110. && !defined(ASIO_NO_EXCEPTIONS)
  111. }
  112. catch (...)
  113. {
  114. pool_->scheduler_.capture_current_exception();
  115. return;
  116. }
  117. #endif // defined(ASIO_HAS_STD_EXCEPTION_PTR)
  118. // && !defined(ASIO_NO_EXCEPTIONS)
  119. }
  120. // Allocate and construct an operation to wrap the function.
  121. typedef detail::executor_op<function_type, Allocator> op;
  122. typename op::ptr p = { detail::addressof(allocator_),
  123. op::ptr::allocate(allocator_), 0 };
  124. p.p = new (p.v) op(ASIO_MOVE_CAST(Function)(f), allocator_);
  125. if ((bits_ & relationship_continuation) != 0)
  126. {
  127. ASIO_HANDLER_CREATION((*pool_, *p.p,
  128. "thread_pool", pool_, 0, "execute(blk=never,rel=cont)"));
  129. }
  130. else
  131. {
  132. ASIO_HANDLER_CREATION((*pool_, *p.p,
  133. "thread_pool", pool_, 0, "execute(blk=never,rel=fork)"));
  134. }
  135. pool_->scheduler_.post_immediate_completion(p.p,
  136. (bits_ & relationship_continuation) != 0);
  137. p.v = p.p = 0;
  138. }
  139. template <typename Allocator, unsigned int Bits>
  140. template <typename Function>
  141. void thread_pool::basic_executor_type<Allocator,
  142. Bits>::do_execute(ASIO_MOVE_ARG(Function) f, true_type) const
  143. {
  144. // Obtain a non-const instance of the function.
  145. detail::non_const_lvalue<Function> f2(f);
  146. // Invoke immediately if we are already inside the thread pool.
  147. if (pool_->scheduler_.can_dispatch())
  148. {
  149. #if !defined(ASIO_NO_EXCEPTIONS)
  150. try
  151. {
  152. #endif // !defined(ASIO_NO_EXCEPTIONS)
  153. detail::fenced_block b(detail::fenced_block::full);
  154. asio_handler_invoke_helpers::invoke(f2.value, f2.value);
  155. return;
  156. #if !defined(ASIO_NO_EXCEPTIONS)
  157. }
  158. catch (...)
  159. {
  160. std::terminate();
  161. }
  162. #endif // !defined(ASIO_NO_EXCEPTIONS)
  163. }
  164. // Construct an operation to wrap the function.
  165. typedef typename decay<Function>::type function_type;
  166. detail::blocking_executor_op<function_type> op(f2.value);
  167. ASIO_HANDLER_CREATION((*pool_, op,
  168. "thread_pool", pool_, 0, "execute(blk=always)"));
  169. pool_->scheduler_.post_immediate_completion(&op, false);
  170. op.wait();
  171. }
  172. template <typename Allocator, unsigned int Bits>
  173. template <typename Function>
  174. void thread_pool::basic_executor_type<Allocator, Bits>::do_bulk_execute(
  175. ASIO_MOVE_ARG(Function) f, std::size_t n, false_type) const
  176. {
  177. typedef typename decay<Function>::type function_type;
  178. typedef detail::bulk_executor_op<function_type, Allocator> op;
  179. // Allocate and construct operations to wrap the function.
  180. detail::op_queue<detail::scheduler_operation> ops;
  181. for (std::size_t i = 0; i < n; ++i)
  182. {
  183. typename op::ptr p = { detail::addressof(allocator_),
  184. op::ptr::allocate(allocator_), 0 };
  185. p.p = new (p.v) op(ASIO_MOVE_CAST(Function)(f), allocator_, i);
  186. ops.push(p.p);
  187. if ((bits_ & relationship_continuation) != 0)
  188. {
  189. ASIO_HANDLER_CREATION((*pool_, *p.p,
  190. "thread_pool", pool_, 0, "bulk_execute(blk=never,rel=cont)"));
  191. }
  192. else
  193. {
  194. ASIO_HANDLER_CREATION((*pool_, *p.p,
  195. "thread_pool", pool_, 0, "bulk)execute(blk=never,rel=fork)"));
  196. }
  197. p.v = p.p = 0;
  198. }
  199. pool_->scheduler_.post_immediate_completions(n,
  200. ops, (bits_ & relationship_continuation) != 0);
  201. }
  202. template <typename Function>
  203. struct thread_pool_always_blocking_function_adapter
  204. {
  205. typename decay<Function>::type* f;
  206. std::size_t n;
  207. void operator()()
  208. {
  209. for (std::size_t i = 0; i < n; ++i)
  210. {
  211. (*f)(i);
  212. }
  213. }
  214. };
  215. template <typename Allocator, unsigned int Bits>
  216. template <typename Function>
  217. void thread_pool::basic_executor_type<Allocator, Bits>::do_bulk_execute(
  218. ASIO_MOVE_ARG(Function) f, std::size_t n, true_type) const
  219. {
  220. // Obtain a non-const instance of the function.
  221. detail::non_const_lvalue<Function> f2(f);
  222. thread_pool_always_blocking_function_adapter<Function>
  223. adapter = { detail::addressof(f2.value), n };
  224. this->do_execute(adapter, true_type());
  225. }
  226. #if !defined(ASIO_NO_TS_EXECUTORS)
  227. template <typename Allocator, unsigned int Bits>
  228. inline thread_pool& thread_pool::basic_executor_type<
  229. Allocator, Bits>::context() const ASIO_NOEXCEPT
  230. {
  231. return *pool_;
  232. }
  233. template <typename Allocator, unsigned int Bits>
  234. inline void thread_pool::basic_executor_type<Allocator,
  235. Bits>::on_work_started() const ASIO_NOEXCEPT
  236. {
  237. pool_->scheduler_.work_started();
  238. }
  239. template <typename Allocator, unsigned int Bits>
  240. inline void thread_pool::basic_executor_type<Allocator,
  241. Bits>::on_work_finished() const ASIO_NOEXCEPT
  242. {
  243. pool_->scheduler_.work_finished();
  244. }
  245. template <typename Allocator, unsigned int Bits>
  246. template <typename Function, typename OtherAllocator>
  247. void thread_pool::basic_executor_type<Allocator, Bits>::dispatch(
  248. ASIO_MOVE_ARG(Function) f, const OtherAllocator& a) const
  249. {
  250. typedef typename decay<Function>::type function_type;
  251. // Invoke immediately if we are already inside the thread pool.
  252. if (pool_->scheduler_.can_dispatch())
  253. {
  254. // Make a local, non-const copy of the function.
  255. function_type tmp(ASIO_MOVE_CAST(Function)(f));
  256. detail::fenced_block b(detail::fenced_block::full);
  257. asio_handler_invoke_helpers::invoke(tmp, tmp);
  258. return;
  259. }
  260. // Allocate and construct an operation to wrap the function.
  261. typedef detail::executor_op<function_type, OtherAllocator> op;
  262. typename op::ptr p = { detail::addressof(a), op::ptr::allocate(a), 0 };
  263. p.p = new (p.v) op(ASIO_MOVE_CAST(Function)(f), a);
  264. ASIO_HANDLER_CREATION((*pool_, *p.p,
  265. "thread_pool", pool_, 0, "dispatch"));
  266. pool_->scheduler_.post_immediate_completion(p.p, false);
  267. p.v = p.p = 0;
  268. }
  269. template <typename Allocator, unsigned int Bits>
  270. template <typename Function, typename OtherAllocator>
  271. void thread_pool::basic_executor_type<Allocator, Bits>::post(
  272. ASIO_MOVE_ARG(Function) f, const OtherAllocator& a) const
  273. {
  274. typedef typename decay<Function>::type function_type;
  275. // Allocate and construct an operation to wrap the function.
  276. typedef detail::executor_op<function_type, OtherAllocator> op;
  277. typename op::ptr p = { detail::addressof(a), op::ptr::allocate(a), 0 };
  278. p.p = new (p.v) op(ASIO_MOVE_CAST(Function)(f), a);
  279. ASIO_HANDLER_CREATION((*pool_, *p.p,
  280. "thread_pool", pool_, 0, "post"));
  281. pool_->scheduler_.post_immediate_completion(p.p, false);
  282. p.v = p.p = 0;
  283. }
  284. template <typename Allocator, unsigned int Bits>
  285. template <typename Function, typename OtherAllocator>
  286. void thread_pool::basic_executor_type<Allocator, Bits>::defer(
  287. ASIO_MOVE_ARG(Function) f, const OtherAllocator& a) const
  288. {
  289. typedef typename decay<Function>::type function_type;
  290. // Allocate and construct an operation to wrap the function.
  291. typedef detail::executor_op<function_type, OtherAllocator> op;
  292. typename op::ptr p = { detail::addressof(a), op::ptr::allocate(a), 0 };
  293. p.p = new (p.v) op(ASIO_MOVE_CAST(Function)(f), a);
  294. ASIO_HANDLER_CREATION((*pool_, *p.p,
  295. "thread_pool", pool_, 0, "defer"));
  296. pool_->scheduler_.post_immediate_completion(p.p, true);
  297. p.v = p.p = 0;
  298. }
  299. #endif // !defined(ASIO_NO_TS_EXECUTORS)
  300. } // namespace asio
  301. #include "asio/detail/pop_options.hpp"
  302. #endif // ASIO_IMPL_THREAD_POOL_HPP