basic_concurrent_channel.hpp 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431
  1. //
  2. // experimental/basic_concurrent_channel.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_EXPERIMENTAL_BASIC_CONCURRENT_CHANNEL_HPP
  11. #define ASIO_EXPERIMENTAL_BASIC_CONCURRENT_CHANNEL_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 "asio/detail/non_const_lvalue.hpp"
  17. #include "asio/detail/mutex.hpp"
  18. #include "asio/execution/executor.hpp"
  19. #include "asio/execution_context.hpp"
  20. #include "asio/experimental/detail/channel_send_functions.hpp"
  21. #include "asio/experimental/detail/channel_service.hpp"
  22. #include "asio/detail/push_options.hpp"
  23. namespace asio {
  24. namespace experimental {
  25. namespace detail {
  26. } // namespace detail
  27. /// A channel for messages.
  28. template <typename Executor, typename Traits, typename... Signatures>
  29. class basic_concurrent_channel
  30. #if !defined(GENERATING_DOCUMENTATION)
  31. : public detail::channel_send_functions<
  32. basic_concurrent_channel<Executor, Traits, Signatures...>,
  33. Executor, Signatures...>
  34. #endif // !defined(GENERATING_DOCUMENTATION)
  35. {
  36. private:
  37. class initiate_async_send;
  38. class initiate_async_receive;
  39. typedef detail::channel_service<asio::detail::mutex> service_type;
  40. typedef typename service_type::template implementation_type<
  41. Traits, Signatures...>::payload_type payload_type;
  42. template <typename... PayloadSignatures,
  43. ASIO_COMPLETION_TOKEN_FOR(PayloadSignatures...) CompletionToken>
  44. auto do_async_receive(detail::channel_payload<PayloadSignatures...>*,
  45. ASIO_MOVE_ARG(CompletionToken) token)
  46. -> decltype(
  47. async_initiate<CompletionToken, PayloadSignatures...>(
  48. declval<initiate_async_receive>(), token))
  49. {
  50. return async_initiate<CompletionToken, PayloadSignatures...>(
  51. initiate_async_receive(this), token);
  52. }
  53. public:
  54. /// The type of the executor associated with the channel.
  55. typedef Executor executor_type;
  56. /// Rebinds the channel type to another executor.
  57. template <typename Executor1>
  58. struct rebind_executor
  59. {
  60. /// The channel type when rebound to the specified executor.
  61. typedef basic_concurrent_channel<Executor1, Traits, Signatures...> other;
  62. };
  63. /// The traits type associated with the channel.
  64. typedef typename Traits::template rebind<Signatures...>::other traits_type;
  65. /// Construct a basic_concurrent_channel.
  66. /**
  67. * This constructor creates and channel.
  68. *
  69. * @param ex The I/O executor that the channel will use, by default, to
  70. * dispatch handlers for any asynchronous operations performed on the channel.
  71. *
  72. * @param max_buffer_size The maximum number of messages that may be buffered
  73. * in the channel.
  74. */
  75. basic_concurrent_channel(const executor_type& ex,
  76. std::size_t max_buffer_size = 0)
  77. : service_(&asio::use_service<service_type>(
  78. basic_concurrent_channel::get_context(ex))),
  79. impl_(),
  80. executor_(ex)
  81. {
  82. service_->construct(impl_, max_buffer_size);
  83. }
  84. /// Construct and open a basic_concurrent_channel.
  85. /**
  86. * This constructor creates and opens a channel.
  87. *
  88. * @param context An execution context which provides the I/O executor that
  89. * the channel will use, by default, to dispatch handlers for any asynchronous
  90. * operations performed on the channel.
  91. *
  92. * @param max_buffer_size The maximum number of messages that may be buffered
  93. * in the channel.
  94. */
  95. template <typename ExecutionContext>
  96. basic_concurrent_channel(ExecutionContext& context,
  97. std::size_t max_buffer_size = 0,
  98. typename constraint<
  99. is_convertible<ExecutionContext&, execution_context&>::value,
  100. defaulted_constraint
  101. >::type = defaulted_constraint())
  102. : service_(&asio::use_service<service_type>(context)),
  103. impl_(),
  104. executor_(context.get_executor())
  105. {
  106. service_->construct(impl_, max_buffer_size);
  107. }
  108. #if defined(ASIO_HAS_MOVE) || defined(GENERATING_DOCUMENTATION)
  109. /// Move-construct a basic_concurrent_channel from another.
  110. /**
  111. * This constructor moves a channel from one object to another.
  112. *
  113. * @param other The other basic_concurrent_channel object from which the move
  114. * will occur.
  115. *
  116. * @note Following the move, the moved-from object is in the same state as if
  117. * constructed using the @c basic_concurrent_channel(const executor_type&)
  118. * constructor.
  119. */
  120. basic_concurrent_channel(basic_concurrent_channel&& other)
  121. : service_(other.service_),
  122. executor_(other.executor_)
  123. {
  124. service_->move_construct(impl_, other.impl_);
  125. }
  126. /// Move-assign a basic_concurrent_channel from another.
  127. /**
  128. * This assignment operator moves a channel from one object to another.
  129. * Cancels any outstanding asynchronous operations associated with the target
  130. * object.
  131. *
  132. * @param other The other basic_concurrent_channel object from which the move
  133. * will occur.
  134. *
  135. * @note Following the move, the moved-from object is in the same state as if
  136. * constructed using the @c basic_concurrent_channel(const executor_type&)
  137. * constructor.
  138. */
  139. basic_concurrent_channel& operator=(basic_concurrent_channel&& other)
  140. {
  141. if (this != &other)
  142. {
  143. service_->move_assign(impl_, *other.service_, other.impl_);
  144. executor_.~executor_type();
  145. new (&executor_) executor_type(other.executor_);
  146. service_ = other.service_;
  147. }
  148. return *this;
  149. }
  150. // All channels have access to each other's implementations.
  151. template <typename, typename, typename...>
  152. friend class basic_concurrent_channel;
  153. /// Move-construct a basic_concurrent_channel from another.
  154. /**
  155. * This constructor moves a channel from one object to another.
  156. *
  157. * @param other The other basic_concurrent_channel object from which the move
  158. * will occur.
  159. *
  160. * @note Following the move, the moved-from object is in the same state as if
  161. * constructed using the @c basic_concurrent_channel(const executor_type&)
  162. * constructor.
  163. */
  164. template <typename Executor1>
  165. basic_concurrent_channel(
  166. basic_concurrent_channel<Executor1, Traits, Signatures...>&& other,
  167. typename constraint<
  168. is_convertible<Executor1, Executor>::value
  169. >::type = 0)
  170. : service_(other.service_),
  171. executor_(other.executor_)
  172. {
  173. service_->move_construct(impl_, *other.service_, other.impl_);
  174. }
  175. /// Move-assign a basic_concurrent_channel from another.
  176. /**
  177. * This assignment operator moves a channel from one object to another.
  178. * Cancels any outstanding asynchronous operations associated with the target
  179. * object.
  180. *
  181. * @param other The other basic_concurrent_channel object from which the move
  182. * will occur.
  183. *
  184. * @note Following the move, the moved-from object is in the same state as if
  185. * constructed using the @c basic_concurrent_channel(const executor_type&)
  186. * constructor.
  187. */
  188. template <typename Executor1>
  189. typename constraint<
  190. is_convertible<Executor1, Executor>::value,
  191. basic_concurrent_channel&
  192. >::type operator=(
  193. basic_concurrent_channel<Executor1, Traits, Signatures...>&& other)
  194. {
  195. if (this != &other)
  196. {
  197. service_->move_assign(impl_, *other.service_, other.impl_);
  198. executor_.~executor_type();
  199. new (&executor_) executor_type(other.executor_);
  200. service_ = other.service_;
  201. }
  202. return *this;
  203. }
  204. #endif // defined(ASIO_HAS_MOVE) || defined(GENERATING_DOCUMENTATION)
  205. /// Destructor.
  206. ~basic_concurrent_channel()
  207. {
  208. service_->destroy(impl_);
  209. }
  210. /// Get the executor associated with the object.
  211. executor_type get_executor() ASIO_NOEXCEPT
  212. {
  213. return executor_;
  214. }
  215. /// Get the capacity of the channel's buffer.
  216. std::size_t capacity() ASIO_NOEXCEPT
  217. {
  218. return service_->capacity(impl_);
  219. }
  220. /// Determine whether the channel is open.
  221. bool is_open() const ASIO_NOEXCEPT
  222. {
  223. return service_->is_open(impl_);
  224. }
  225. /// Reset the channel to its initial state.
  226. void reset()
  227. {
  228. service_->reset(impl_);
  229. }
  230. /// Close the channel.
  231. void close()
  232. {
  233. service_->close(impl_);
  234. }
  235. /// Cancel all asynchronous operations waiting on the channel.
  236. /**
  237. * All outstanding send operations will complete with the error
  238. * @c asio::experimental::error::channel_canceld. Outstanding receive
  239. * operations complete with the result as determined by the channel traits.
  240. */
  241. void cancel()
  242. {
  243. service_->cancel(impl_);
  244. }
  245. /// Determine whether a message can be received without blocking.
  246. bool ready() const ASIO_NOEXCEPT
  247. {
  248. return service_->ready(impl_);
  249. }
  250. #if defined(GENERATING_DOCUMENTATION)
  251. /// Try to send a message without blocking.
  252. /**
  253. * Fails if the buffer is full and there are no waiting receive operations.
  254. *
  255. * @returns @c true on success, @c false on failure.
  256. */
  257. template <typename... Args>
  258. bool try_send(ASIO_MOVE_ARG(Args)... args);
  259. /// Try to send a number of messages without blocking.
  260. /**
  261. * @returns The number of messages that were sent.
  262. */
  263. template <typename... Args>
  264. std::size_t try_send_n(std::size_t count, ASIO_MOVE_ARG(Args)... args);
  265. /// Asynchronously send a message.
  266. template <typename... Args,
  267. ASIO_COMPLETION_TOKEN_FOR(void (asio::error_code))
  268. CompletionToken ASIO_DEFAULT_COMPLETION_TOKEN_TYPE(executor_type)>
  269. auto async_send(ASIO_MOVE_ARG(Args)... args,
  270. ASIO_MOVE_ARG(CompletionToken) token);
  271. #endif // defined(GENERATING_DOCUMENTATION)
  272. /// Try to receive a message without blocking.
  273. /**
  274. * Fails if the buffer is full and there are no waiting receive operations.
  275. *
  276. * @returns @c true on success, @c false on failure.
  277. */
  278. template <typename Handler>
  279. bool try_receive(ASIO_MOVE_ARG(Handler) handler)
  280. {
  281. return service_->try_receive(impl_, ASIO_MOVE_CAST(Handler)(handler));
  282. }
  283. /// Asynchronously receive a message.
  284. template <typename CompletionToken
  285. ASIO_DEFAULT_COMPLETION_TOKEN_TYPE(executor_type)>
  286. auto async_receive(
  287. ASIO_MOVE_ARG(CompletionToken) token
  288. ASIO_DEFAULT_COMPLETION_TOKEN(Executor))
  289. #if !defined(GENERATING_DOCUMENTATION)
  290. -> decltype(
  291. this->do_async_receive(static_cast<payload_type*>(0),
  292. ASIO_MOVE_CAST(CompletionToken)(token)))
  293. #endif // !defined(GENERATING_DOCUMENTATION)
  294. {
  295. return this->do_async_receive(static_cast<payload_type*>(0),
  296. ASIO_MOVE_CAST(CompletionToken)(token));
  297. }
  298. private:
  299. // Disallow copying and assignment.
  300. basic_concurrent_channel(
  301. const basic_concurrent_channel&) ASIO_DELETED;
  302. basic_concurrent_channel& operator=(
  303. const basic_concurrent_channel&) ASIO_DELETED;
  304. template <typename, typename, typename...>
  305. friend class detail::channel_send_functions;
  306. // Helper function to get an executor's context.
  307. template <typename T>
  308. static execution_context& get_context(const T& t,
  309. typename enable_if<execution::is_executor<T>::value>::type* = 0)
  310. {
  311. return asio::query(t, execution::context);
  312. }
  313. // Helper function to get an executor's context.
  314. template <typename T>
  315. static execution_context& get_context(const T& t,
  316. typename enable_if<!execution::is_executor<T>::value>::type* = 0)
  317. {
  318. return t.context();
  319. }
  320. class initiate_async_send
  321. {
  322. public:
  323. typedef Executor executor_type;
  324. explicit initiate_async_send(basic_concurrent_channel* self)
  325. : self_(self)
  326. {
  327. }
  328. executor_type get_executor() const ASIO_NOEXCEPT
  329. {
  330. return self_->get_executor();
  331. }
  332. template <typename SendHandler>
  333. void operator()(ASIO_MOVE_ARG(SendHandler) handler,
  334. ASIO_MOVE_ARG(payload_type) payload) const
  335. {
  336. asio::detail::non_const_lvalue<SendHandler> handler2(handler);
  337. self_->service_->async_send(self_->impl_,
  338. ASIO_MOVE_CAST(payload_type)(payload),
  339. handler2.value, self_->get_executor());
  340. }
  341. private:
  342. basic_concurrent_channel* self_;
  343. };
  344. class initiate_async_receive
  345. {
  346. public:
  347. typedef Executor executor_type;
  348. explicit initiate_async_receive(basic_concurrent_channel* self)
  349. : self_(self)
  350. {
  351. }
  352. executor_type get_executor() const ASIO_NOEXCEPT
  353. {
  354. return self_->get_executor();
  355. }
  356. template <typename ReceiveHandler>
  357. void operator()(ASIO_MOVE_ARG(ReceiveHandler) handler) const
  358. {
  359. asio::detail::non_const_lvalue<ReceiveHandler> handler2(handler);
  360. self_->service_->async_receive(self_->impl_,
  361. handler2.value, self_->get_executor());
  362. }
  363. private:
  364. basic_concurrent_channel* self_;
  365. };
  366. // The service associated with the I/O object.
  367. service_type* service_;
  368. // The underlying implementation of the I/O object.
  369. typename service_type::template implementation_type<
  370. Traits, Signatures...> impl_;
  371. // The associated executor.
  372. Executor executor_;
  373. };
  374. } // namespace experimental
  375. } // namespace asio
  376. #include "asio/detail/pop_options.hpp"
  377. #endif // ASIO_EXPERIMENTAL_BASIC_CONCURRENT_CHANNEL_HPP