io_uring_socket_recv_op.hpp 6.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201
  1. //
  2. // detail/io_uring_socket_recv_op.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_DETAIL_IO_URING_SOCKET_RECV_OP_HPP
  11. #define ASIO_DETAIL_IO_URING_SOCKET_RECV_OP_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. #if defined(ASIO_HAS_IO_URING)
  17. #include "asio/detail/bind_handler.hpp"
  18. #include "asio/detail/buffer_sequence_adapter.hpp"
  19. #include "asio/detail/socket_ops.hpp"
  20. #include "asio/detail/fenced_block.hpp"
  21. #include "asio/detail/handler_work.hpp"
  22. #include "asio/detail/io_uring_operation.hpp"
  23. #include "asio/detail/memory.hpp"
  24. #include "asio/detail/push_options.hpp"
  25. namespace asio {
  26. namespace detail {
  27. template <typename MutableBufferSequence>
  28. class io_uring_socket_recv_op_base : public io_uring_operation
  29. {
  30. public:
  31. io_uring_socket_recv_op_base(const asio::error_code& success_ec,
  32. socket_type socket, socket_ops::state_type state,
  33. const MutableBufferSequence& buffers,
  34. socket_base::message_flags flags, func_type complete_func)
  35. : io_uring_operation(success_ec,
  36. &io_uring_socket_recv_op_base::do_prepare,
  37. &io_uring_socket_recv_op_base::do_perform, complete_func),
  38. socket_(socket),
  39. state_(state),
  40. buffers_(buffers),
  41. flags_(flags),
  42. bufs_(buffers),
  43. msghdr_()
  44. {
  45. msghdr_.msg_iov = bufs_.buffers();
  46. msghdr_.msg_iovlen = static_cast<int>(bufs_.count());
  47. }
  48. static void do_prepare(io_uring_operation* base, ::io_uring_sqe* sqe)
  49. {
  50. io_uring_socket_recv_op_base* o(
  51. static_cast<io_uring_socket_recv_op_base*>(base));
  52. if ((o->state_ & socket_ops::internal_non_blocking) != 0)
  53. {
  54. bool except_op = (o->flags_ & socket_base::message_out_of_band) != 0;
  55. ::io_uring_prep_poll_add(sqe, o->socket_, except_op ? POLLPRI : POLLIN);
  56. }
  57. else if (o->bufs_.is_single_buffer
  58. && o->bufs_.is_registered_buffer && o->flags_ == 0)
  59. {
  60. ::io_uring_prep_read_fixed(sqe, o->socket_,
  61. o->bufs_.buffers()->iov_base, o->bufs_.buffers()->iov_len,
  62. 0, o->bufs_.registered_id().native_handle());
  63. }
  64. else
  65. {
  66. ::io_uring_prep_recvmsg(sqe, o->socket_, &o->msghdr_, o->flags_);
  67. }
  68. }
  69. static bool do_perform(io_uring_operation* base, bool after_completion)
  70. {
  71. io_uring_socket_recv_op_base* o(
  72. static_cast<io_uring_socket_recv_op_base*>(base));
  73. if ((o->state_ & socket_ops::internal_non_blocking) != 0)
  74. {
  75. bool except_op = (o->flags_ & socket_base::message_out_of_band) != 0;
  76. if (after_completion || !except_op)
  77. {
  78. if (o->bufs_.is_single_buffer)
  79. {
  80. return socket_ops::non_blocking_recv1(o->socket_,
  81. o->bufs_.first(o->buffers_).data(),
  82. o->bufs_.first(o->buffers_).size(), o->flags_,
  83. (o->state_ & socket_ops::stream_oriented) != 0,
  84. o->ec_, o->bytes_transferred_);
  85. }
  86. else
  87. {
  88. return socket_ops::non_blocking_recv(o->socket_,
  89. o->bufs_.buffers(), o->bufs_.count(), o->flags_,
  90. (o->state_ & socket_ops::stream_oriented) != 0,
  91. o->ec_, o->bytes_transferred_);
  92. }
  93. }
  94. }
  95. else if (after_completion)
  96. {
  97. if (!o->ec_ && o->bytes_transferred_ == 0)
  98. if ((o->state_ & socket_ops::stream_oriented) != 0)
  99. o->ec_ = asio::error::eof;
  100. }
  101. if (o->ec_ && o->ec_ == asio::error::would_block)
  102. {
  103. o->state_ |= socket_ops::internal_non_blocking;
  104. return false;
  105. }
  106. return after_completion;
  107. }
  108. private:
  109. socket_type socket_;
  110. socket_ops::state_type state_;
  111. MutableBufferSequence buffers_;
  112. socket_base::message_flags flags_;
  113. buffer_sequence_adapter<asio::mutable_buffer,
  114. MutableBufferSequence> bufs_;
  115. msghdr msghdr_;
  116. };
  117. template <typename MutableBufferSequence, typename Handler, typename IoExecutor>
  118. class io_uring_socket_recv_op
  119. : public io_uring_socket_recv_op_base<MutableBufferSequence>
  120. {
  121. public:
  122. ASIO_DEFINE_HANDLER_PTR(io_uring_socket_recv_op);
  123. io_uring_socket_recv_op(const asio::error_code& success_ec,
  124. int socket, socket_ops::state_type state,
  125. const MutableBufferSequence& buffers, socket_base::message_flags flags,
  126. Handler& handler, const IoExecutor& io_ex)
  127. : io_uring_socket_recv_op_base<MutableBufferSequence>(success_ec,
  128. socket, state, buffers, flags, &io_uring_socket_recv_op::do_complete),
  129. handler_(ASIO_MOVE_CAST(Handler)(handler)),
  130. work_(handler_, io_ex)
  131. {
  132. }
  133. static void do_complete(void* owner, operation* base,
  134. const asio::error_code& /*ec*/,
  135. std::size_t /*bytes_transferred*/)
  136. {
  137. // Take ownership of the handler object.
  138. io_uring_socket_recv_op* o
  139. (static_cast<io_uring_socket_recv_op*>(base));
  140. ptr p = { asio::detail::addressof(o->handler_), o, o };
  141. ASIO_HANDLER_COMPLETION((*o));
  142. // Take ownership of the operation's outstanding work.
  143. handler_work<Handler, IoExecutor> w(
  144. ASIO_MOVE_CAST2(handler_work<Handler, IoExecutor>)(
  145. o->work_));
  146. // Make a copy of the handler so that the memory can be deallocated before
  147. // the upcall is made. Even if we're not about to make an upcall, a
  148. // sub-object of the handler may be the true owner of the memory associated
  149. // with the handler. Consequently, a local copy of the handler is required
  150. // to ensure that any owning sub-object remains valid until after we have
  151. // deallocated the memory here.
  152. detail::binder2<Handler, asio::error_code, std::size_t>
  153. handler(o->handler_, o->ec_, o->bytes_transferred_);
  154. p.h = asio::detail::addressof(handler.handler_);
  155. p.reset();
  156. // Make the upcall if required.
  157. if (owner)
  158. {
  159. fenced_block b(fenced_block::half);
  160. ASIO_HANDLER_INVOCATION_BEGIN((handler.arg1_, handler.arg2_));
  161. w.complete(handler, handler.handler_);
  162. ASIO_HANDLER_INVOCATION_END;
  163. }
  164. }
  165. private:
  166. Handler handler_;
  167. handler_work<Handler, IoExecutor> work_;
  168. };
  169. } // namespace detail
  170. } // namespace asio
  171. #include "asio/detail/pop_options.hpp"
  172. #endif // defined(ASIO_HAS_IO_URING)
  173. #endif // ASIO_DETAIL_IO_URING_SOCKET_RECV_OP_HPP