channel_service.hpp 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612
  1. //
  2. // experimental/detail/impl/channel_service.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_DETAIL_IMPL_CHANNEL_SERVICE_HPP
  11. #define ASIO_EXPERIMENTAL_DETAIL_IMPL_CHANNEL_SERVICE_HPP
  12. #if defined(_MSC_VER) && (_MSC_VER >= 1200)
  13. # pragma once
  14. #endif // defined(_MSC_VER) && (_MSC_VER >= 1200)
  15. #include "asio/detail/push_options.hpp"
  16. namespace asio {
  17. namespace experimental {
  18. namespace detail {
  19. template <typename Mutex>
  20. inline channel_service<Mutex>::channel_service(execution_context& ctx)
  21. : asio::detail::execution_context_service_base<channel_service>(ctx),
  22. mutex_(),
  23. impl_list_(0)
  24. {
  25. }
  26. template <typename Mutex>
  27. inline void channel_service<Mutex>::shutdown()
  28. {
  29. // Abandon all pending operations.
  30. asio::detail::op_queue<channel_operation> ops;
  31. asio::detail::mutex::scoped_lock lock(mutex_);
  32. base_implementation_type* impl = impl_list_;
  33. while (impl)
  34. {
  35. ops.push(impl->waiters_);
  36. impl = impl->next_;
  37. }
  38. }
  39. template <typename Mutex>
  40. inline void channel_service<Mutex>::construct(
  41. channel_service<Mutex>::base_implementation_type& impl,
  42. std::size_t max_buffer_size)
  43. {
  44. impl.max_buffer_size_ = max_buffer_size;
  45. impl.receive_state_ = block;
  46. impl.send_state_ = max_buffer_size ? buffer : block;
  47. // Insert implementation into linked list of all implementations.
  48. asio::detail::mutex::scoped_lock lock(mutex_);
  49. impl.next_ = impl_list_;
  50. impl.prev_ = 0;
  51. if (impl_list_)
  52. impl_list_->prev_ = &impl;
  53. impl_list_ = &impl;
  54. }
  55. template <typename Mutex>
  56. template <typename Traits, typename... Signatures>
  57. void channel_service<Mutex>::destroy(
  58. channel_service<Mutex>::implementation_type<Traits, Signatures...>& impl)
  59. {
  60. cancel(impl);
  61. base_destroy(impl);
  62. }
  63. template <typename Mutex>
  64. template <typename Traits, typename... Signatures>
  65. void channel_service<Mutex>::move_construct(
  66. channel_service<Mutex>::implementation_type<Traits, Signatures...>& impl,
  67. channel_service<Mutex>::implementation_type<
  68. Traits, Signatures...>& other_impl)
  69. {
  70. impl.max_buffer_size_ = other_impl.max_buffer_size_;
  71. impl.receive_state_ = other_impl.receive_state_;
  72. other_impl.receive_state_ = block;
  73. impl.send_state_ = other_impl.send_state_;
  74. other_impl.send_state_ = other_impl.max_buffer_size_ ? buffer : block;
  75. impl.buffer_move_from(other_impl);
  76. // Insert implementation into linked list of all implementations.
  77. asio::detail::mutex::scoped_lock lock(mutex_);
  78. impl.next_ = impl_list_;
  79. impl.prev_ = 0;
  80. if (impl_list_)
  81. impl_list_->prev_ = &impl;
  82. impl_list_ = &impl;
  83. }
  84. template <typename Mutex>
  85. template <typename Traits, typename... Signatures>
  86. void channel_service<Mutex>::move_assign(
  87. channel_service<Mutex>::implementation_type<Traits, Signatures...>& impl,
  88. channel_service& other_service,
  89. channel_service<Mutex>::implementation_type<
  90. Traits, Signatures...>& other_impl)
  91. {
  92. cancel(impl);
  93. if (this != &other_service)
  94. {
  95. // Remove implementation from linked list of all implementations.
  96. asio::detail::mutex::scoped_lock lock(mutex_);
  97. if (impl_list_ == &impl)
  98. impl_list_ = impl.next_;
  99. if (impl.prev_)
  100. impl.prev_->next_ = impl.next_;
  101. if (impl.next_)
  102. impl.next_->prev_= impl.prev_;
  103. impl.next_ = 0;
  104. impl.prev_ = 0;
  105. }
  106. impl.max_buffer_size_ = other_impl.max_buffer_size_;
  107. impl.receive_state_ = other_impl.receive_state_;
  108. other_impl.receive_state_ = block;
  109. impl.send_state_ = other_impl.send_state_;
  110. other_impl.send_state_ = other_impl.max_buffer_size_ ? buffer : block;
  111. impl.buffer_move_from(other_impl);
  112. if (this != &other_service)
  113. {
  114. // Insert implementation into linked list of all implementations.
  115. asio::detail::mutex::scoped_lock lock(other_service.mutex_);
  116. impl.next_ = other_service.impl_list_;
  117. impl.prev_ = 0;
  118. if (other_service.impl_list_)
  119. other_service.impl_list_->prev_ = &impl;
  120. other_service.impl_list_ = &impl;
  121. }
  122. }
  123. template <typename Mutex>
  124. inline void channel_service<Mutex>::base_destroy(
  125. channel_service<Mutex>::base_implementation_type& impl)
  126. {
  127. // Remove implementation from linked list of all implementations.
  128. asio::detail::mutex::scoped_lock lock(mutex_);
  129. if (impl_list_ == &impl)
  130. impl_list_ = impl.next_;
  131. if (impl.prev_)
  132. impl.prev_->next_ = impl.next_;
  133. if (impl.next_)
  134. impl.next_->prev_= impl.prev_;
  135. impl.next_ = 0;
  136. impl.prev_ = 0;
  137. }
  138. template <typename Mutex>
  139. inline std::size_t channel_service<Mutex>::capacity(
  140. const channel_service<Mutex>::base_implementation_type& impl)
  141. const ASIO_NOEXCEPT
  142. {
  143. typename Mutex::scoped_lock lock(impl.mutex_);
  144. return impl.max_buffer_size_;
  145. }
  146. template <typename Mutex>
  147. inline bool channel_service<Mutex>::is_open(
  148. const channel_service<Mutex>::base_implementation_type& impl)
  149. const ASIO_NOEXCEPT
  150. {
  151. typename Mutex::scoped_lock lock(impl.mutex_);
  152. return impl.send_state_ != closed;
  153. }
  154. template <typename Mutex>
  155. template <typename Traits, typename... Signatures>
  156. void channel_service<Mutex>::reset(
  157. channel_service<Mutex>::implementation_type<Traits, Signatures...>& impl)
  158. {
  159. cancel(impl);
  160. typename Mutex::scoped_lock lock(impl.mutex_);
  161. if (impl.receive_state_ == closed)
  162. impl.receive_state_ = block;
  163. if (impl.send_state_ == closed)
  164. impl.send_state_ = impl.max_buffer_size_ ? buffer : block;
  165. impl.buffer_clear();
  166. }
  167. template <typename Mutex>
  168. template <typename Traits, typename... Signatures>
  169. void channel_service<Mutex>::close(
  170. channel_service<Mutex>::implementation_type<Traits, Signatures...>& impl)
  171. {
  172. typedef typename implementation_type<Traits,
  173. Signatures...>::traits_type traits_type;
  174. typedef typename implementation_type<Traits,
  175. Signatures...>::payload_type payload_type;
  176. typename Mutex::scoped_lock lock(impl.mutex_);
  177. if (impl.receive_state_ == block)
  178. {
  179. while (channel_operation* op = impl.waiters_.front())
  180. {
  181. impl.waiters_.pop();
  182. traits_type::invoke_receive_closed(
  183. complete_receive<payload_type,
  184. typename traits_type::receive_closed_signature>(
  185. static_cast<channel_receive<payload_type>*>(op)));
  186. }
  187. }
  188. impl.send_state_ = closed;
  189. impl.receive_state_ = closed;
  190. }
  191. template <typename Mutex>
  192. template <typename Traits, typename... Signatures>
  193. void channel_service<Mutex>::cancel(
  194. channel_service<Mutex>::implementation_type<Traits, Signatures...>& impl)
  195. {
  196. typedef typename implementation_type<Traits,
  197. Signatures...>::traits_type traits_type;
  198. typedef typename implementation_type<Traits,
  199. Signatures...>::payload_type payload_type;
  200. typename Mutex::scoped_lock lock(impl.mutex_);
  201. while (channel_operation* op = impl.waiters_.front())
  202. {
  203. if (impl.send_state_ == block)
  204. {
  205. impl.waiters_.pop();
  206. static_cast<channel_send<payload_type>*>(op)->cancel();
  207. }
  208. else
  209. {
  210. impl.waiters_.pop();
  211. traits_type::invoke_receive_cancelled(
  212. complete_receive<payload_type,
  213. typename traits_type::receive_cancelled_signature>(
  214. static_cast<channel_receive<payload_type>*>(op)));
  215. }
  216. }
  217. if (impl.receive_state_ == waiter)
  218. impl.receive_state_ = block;
  219. if (impl.send_state_ == waiter)
  220. impl.send_state_ = block;
  221. }
  222. template <typename Mutex>
  223. template <typename Traits, typename... Signatures>
  224. void channel_service<Mutex>::cancel_by_key(
  225. channel_service<Mutex>::implementation_type<Traits, Signatures...>& impl,
  226. void* cancellation_key)
  227. {
  228. typedef typename implementation_type<Traits,
  229. Signatures...>::traits_type traits_type;
  230. typedef typename implementation_type<Traits,
  231. Signatures...>::payload_type payload_type;
  232. typename Mutex::scoped_lock lock(impl.mutex_);
  233. asio::detail::op_queue<channel_operation> other_ops;
  234. while (channel_operation* op = impl.waiters_.front())
  235. {
  236. if (op->cancellation_key_ == cancellation_key)
  237. {
  238. if (impl.send_state_ == block)
  239. {
  240. impl.waiters_.pop();
  241. static_cast<channel_send<payload_type>*>(op)->cancel();
  242. }
  243. else
  244. {
  245. impl.waiters_.pop();
  246. traits_type::invoke_receive_cancelled(
  247. complete_receive<payload_type,
  248. typename traits_type::receive_cancelled_signature>(
  249. static_cast<channel_receive<payload_type>*>(op)));
  250. }
  251. }
  252. else
  253. {
  254. impl.waiters_.pop();
  255. other_ops.push(op);
  256. }
  257. }
  258. impl.waiters_.push(other_ops);
  259. if (impl.waiters_.empty())
  260. {
  261. if (impl.receive_state_ == waiter)
  262. impl.receive_state_ = block;
  263. if (impl.send_state_ == waiter)
  264. impl.send_state_ = block;
  265. }
  266. }
  267. template <typename Mutex>
  268. inline bool channel_service<Mutex>::ready(
  269. const channel_service<Mutex>::base_implementation_type& impl)
  270. const ASIO_NOEXCEPT
  271. {
  272. typename Mutex::scoped_lock lock(impl.mutex_);
  273. return impl.receive_state_ != block;
  274. }
  275. template <typename Mutex>
  276. template <typename Message, typename Traits,
  277. typename... Signatures, typename... Args>
  278. bool channel_service<Mutex>::try_send(
  279. channel_service<Mutex>::implementation_type<Traits, Signatures...>& impl,
  280. ASIO_MOVE_ARG(Args)... args)
  281. {
  282. typedef typename implementation_type<Traits,
  283. Signatures...>::payload_type payload_type;
  284. typename Mutex::scoped_lock lock(impl.mutex_);
  285. switch (impl.send_state_)
  286. {
  287. case block:
  288. {
  289. return false;
  290. }
  291. case buffer:
  292. {
  293. impl.buffer_push(Message(0, ASIO_MOVE_CAST(Args)(args)...));
  294. impl.receive_state_ = buffer;
  295. if (impl.buffer_size() == impl.max_buffer_size_)
  296. impl.send_state_ = block;
  297. return true;
  298. }
  299. case waiter:
  300. {
  301. payload_type payload(Message(0, ASIO_MOVE_CAST(Args)(args)...));
  302. channel_receive<payload_type>* receive_op =
  303. static_cast<channel_receive<payload_type>*>(impl.waiters_.front());
  304. impl.waiters_.pop();
  305. receive_op->complete(ASIO_MOVE_CAST(payload_type)(payload));
  306. if (impl.waiters_.empty())
  307. impl.send_state_ = impl.max_buffer_size_ ? buffer : block;
  308. return true;
  309. }
  310. case closed:
  311. default:
  312. {
  313. return false;
  314. }
  315. }
  316. }
  317. template <typename Mutex>
  318. template <typename Message, typename Traits,
  319. typename... Signatures, typename... Args>
  320. std::size_t channel_service<Mutex>::try_send_n(
  321. channel_service<Mutex>::implementation_type<Traits, Signatures...>& impl,
  322. std::size_t count, ASIO_MOVE_ARG(Args)... args)
  323. {
  324. typedef typename implementation_type<Traits,
  325. Signatures...>::payload_type payload_type;
  326. typename Mutex::scoped_lock lock(impl.mutex_);
  327. if (count == 0)
  328. return 0;
  329. switch (impl.send_state_)
  330. {
  331. case block:
  332. return 0;
  333. case buffer:
  334. case waiter:
  335. break;
  336. case closed:
  337. default:
  338. return 0;
  339. }
  340. payload_type payload(Message(0, ASIO_MOVE_CAST(Args)(args)...));
  341. for (std::size_t i = 0; i < count; ++i)
  342. {
  343. switch (impl.send_state_)
  344. {
  345. case block:
  346. {
  347. return i;
  348. }
  349. case buffer:
  350. {
  351. i += impl.buffer_push_n(count - i,
  352. ASIO_MOVE_CAST(payload_type)(payload));
  353. impl.receive_state_ = buffer;
  354. if (impl.buffer_size() == impl.max_buffer_size_)
  355. impl.send_state_ = block;
  356. return i;
  357. }
  358. case waiter:
  359. {
  360. channel_receive<payload_type>* receive_op =
  361. static_cast<channel_receive<payload_type>*>(impl.waiters_.front());
  362. impl.waiters_.pop();
  363. receive_op->complete(payload);
  364. if (impl.waiters_.empty())
  365. impl.send_state_ = impl.max_buffer_size_ ? buffer : block;
  366. break;
  367. }
  368. case closed:
  369. default:
  370. {
  371. return i;
  372. }
  373. }
  374. }
  375. return count;
  376. }
  377. template <typename Mutex>
  378. template <typename Traits, typename... Signatures>
  379. void channel_service<Mutex>::start_send_op(
  380. channel_service<Mutex>::implementation_type<Traits, Signatures...>& impl,
  381. channel_send<typename implementation_type<
  382. Traits, Signatures...>::payload_type>* send_op)
  383. {
  384. typedef typename implementation_type<Traits,
  385. Signatures...>::payload_type payload_type;
  386. typename Mutex::scoped_lock lock(impl.mutex_);
  387. switch (impl.send_state_)
  388. {
  389. case block:
  390. {
  391. impl.waiters_.push(send_op);
  392. if (impl.receive_state_ == block)
  393. impl.receive_state_ = waiter;
  394. return;
  395. }
  396. case buffer:
  397. {
  398. impl.buffer_push(send_op->get_payload());
  399. impl.receive_state_ = buffer;
  400. if (impl.buffer_size() == impl.max_buffer_size_)
  401. impl.send_state_ = block;
  402. send_op->complete();
  403. break;
  404. }
  405. case waiter:
  406. {
  407. channel_receive<payload_type>* receive_op =
  408. static_cast<channel_receive<payload_type>*>(impl.waiters_.front());
  409. impl.waiters_.pop();
  410. receive_op->complete(send_op->get_payload());
  411. if (impl.waiters_.empty())
  412. impl.send_state_ = impl.max_buffer_size_ ? buffer : block;
  413. send_op->complete();
  414. break;
  415. }
  416. case closed:
  417. default:
  418. {
  419. send_op->close();
  420. break;
  421. }
  422. }
  423. }
  424. template <typename Mutex>
  425. template <typename Traits, typename... Signatures, typename Handler>
  426. bool channel_service<Mutex>::try_receive(
  427. channel_service<Mutex>::implementation_type<Traits, Signatures...>& impl,
  428. ASIO_MOVE_ARG(Handler) handler)
  429. {
  430. typedef typename implementation_type<Traits,
  431. Signatures...>::payload_type payload_type;
  432. typename Mutex::scoped_lock lock(impl.mutex_);
  433. switch (impl.receive_state_)
  434. {
  435. case block:
  436. {
  437. return false;
  438. }
  439. case buffer:
  440. {
  441. payload_type payload(impl.buffer_front());
  442. if (channel_send<payload_type>* send_op =
  443. static_cast<channel_send<payload_type>*>(impl.waiters_.front()))
  444. {
  445. impl.buffer_pop();
  446. impl.buffer_push(send_op->get_payload());
  447. impl.waiters_.pop();
  448. send_op->complete();
  449. }
  450. else
  451. {
  452. impl.buffer_pop();
  453. if (impl.buffer_size() == 0)
  454. impl.receive_state_ = (impl.send_state_ == closed) ? closed : block;
  455. impl.send_state_ = (impl.send_state_ == closed) ? closed : buffer;
  456. }
  457. lock.unlock();
  458. asio::detail::non_const_lvalue<Handler> handler2(handler);
  459. channel_handler<payload_type, typename decay<Handler>::type>(
  460. ASIO_MOVE_CAST(payload_type)(payload), handler2.value)();
  461. return true;
  462. }
  463. case waiter:
  464. {
  465. channel_send<payload_type>* send_op =
  466. static_cast<channel_send<payload_type>*>(impl.waiters_.front());
  467. payload_type payload = send_op->get_payload();
  468. impl.waiters_.pop();
  469. send_op->complete();
  470. if (impl.waiters_.front() == 0)
  471. impl.receive_state_ = (impl.send_state_ == closed) ? closed : block;
  472. lock.unlock();
  473. asio::detail::non_const_lvalue<Handler> handler2(handler);
  474. channel_handler<payload_type, typename decay<Handler>::type>(
  475. ASIO_MOVE_CAST(payload_type)(payload), handler2.value)();
  476. return true;
  477. }
  478. case closed:
  479. default:
  480. {
  481. return false;
  482. }
  483. }
  484. }
  485. template <typename Mutex>
  486. template <typename Traits, typename... Signatures>
  487. void channel_service<Mutex>::start_receive_op(
  488. channel_service<Mutex>::implementation_type<Traits, Signatures...>& impl,
  489. channel_receive<typename implementation_type<
  490. Traits, Signatures...>::payload_type>* receive_op)
  491. {
  492. typedef typename implementation_type<Traits,
  493. Signatures...>::traits_type traits_type;
  494. typedef typename implementation_type<Traits,
  495. Signatures...>::payload_type payload_type;
  496. typename Mutex::scoped_lock lock(impl.mutex_);
  497. switch (impl.receive_state_)
  498. {
  499. case block:
  500. {
  501. impl.waiters_.push(receive_op);
  502. if (impl.send_state_ != closed)
  503. impl.send_state_ = waiter;
  504. return;
  505. }
  506. case buffer:
  507. {
  508. receive_op->complete(impl.buffer_front());
  509. if (channel_send<payload_type>* send_op =
  510. static_cast<channel_send<payload_type>*>(impl.waiters_.front()))
  511. {
  512. impl.buffer_pop();
  513. impl.buffer_push(send_op->get_payload());
  514. impl.waiters_.pop();
  515. send_op->complete();
  516. }
  517. else
  518. {
  519. impl.buffer_pop();
  520. if (impl.buffer_size() == 0)
  521. impl.receive_state_ = (impl.send_state_ == closed) ? closed : block;
  522. impl.send_state_ = (impl.send_state_ == closed) ? closed : buffer;
  523. }
  524. break;
  525. }
  526. case waiter:
  527. {
  528. channel_send<payload_type>* send_op =
  529. static_cast<channel_send<payload_type>*>(impl.waiters_.front());
  530. payload_type payload = send_op->get_payload();
  531. impl.waiters_.pop();
  532. send_op->complete();
  533. receive_op->complete(ASIO_MOVE_CAST(payload_type)(payload));
  534. if (impl.waiters_.front() == 0)
  535. impl.receive_state_ = (impl.send_state_ == closed) ? closed : block;
  536. break;
  537. }
  538. case closed:
  539. default:
  540. {
  541. traits_type::invoke_receive_closed(
  542. complete_receive<payload_type,
  543. typename traits_type::receive_closed_signature>(receive_op));
  544. break;
  545. }
  546. }
  547. }
  548. } // namespace detail
  549. } // namespace experimental
  550. } // namespace asio
  551. #include "asio/detail/pop_options.hpp"
  552. #endif // ASIO_EXPERIMENTAL_DETAIL_IMPL_CHANNEL_SERVICE_HPP