LCOV - code coverage report
Current view: top level - corosio/native/detail/posix - posix_random_access_file_service.hpp (source / functions) Coverage Total Hit Missed
Test: coverage_remapped.info Lines: 99.3 % 143 142 1
Test Date: 2026-10-08 18:13:32 Functions: 100.0 % 15 15

           TLA  Line data    Source code
       1                 : //
       2                 : // Copyright (c) 2026 Michael Vandeberg
       3                 : //
       4                 : // Distributed under the Boost Software License, Version 1.0. (See accompanying
       5                 : // file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
       6                 : //
       7                 : // Official repository: https://github.com/cppalliance/corosio
       8                 : //
       9                 : 
      10                 : #ifndef BOOST_COROSIO_NATIVE_DETAIL_POSIX_POSIX_RANDOM_ACCESS_FILE_SERVICE_HPP
      11                 : #define BOOST_COROSIO_NATIVE_DETAIL_POSIX_POSIX_RANDOM_ACCESS_FILE_SERVICE_HPP
      12                 : 
      13                 : #include <boost/corosio/detail/platform.hpp>
      14                 : 
      15                 : #if BOOST_COROSIO_POSIX
      16                 : 
      17                 : #include <boost/corosio/native/detail/posix/posix_random_access_file.hpp>
      18                 : #include <boost/corosio/native/detail/reactor/reactor_scheduler.hpp>
      19                 : #include <boost/corosio/detail/random_access_file_service.hpp>
      20                 : #include <boost/corosio/detail/object_pool.hpp>
      21                 : #include <boost/corosio/detail/object_ref.hpp>
      22                 : #include <boost/corosio/detail/thread_pool.hpp>
      23                 : 
      24                 : #include <limits>
      25                 : #include <vector>
      26                 : 
      27                 : namespace boost::corosio::detail {
      28                 : 
      29                 : /** Random-access file service for POSIX backends. */
      30                 : class BOOST_COROSIO_DECL posix_random_access_file_service final
      31                 :     : public random_access_file_service
      32                 : {
      33                 :     friend class posix_random_access_file;
      34                 : 
      35                 : public:
      36 HIT         164 :     explicit posix_random_access_file_service(capy::execution_context& ctx)
      37             492 :         : sched_(&get_scheduler(ctx))
      38             164 :         , pool_(ctx)
      39                 :     {
      40             164 :     }
      41                 : 
      42             328 :     ~posix_random_access_file_service() override = default;
      43                 : 
      44                 :     posix_random_access_file_service(posix_random_access_file_service const&) =
      45                 :         delete;
      46                 :     posix_random_access_file_service&
      47                 :     operator=(posix_random_access_file_service const&) = delete;
      48                 : 
      49             219 :     io_object::implementation* construct() override
      50                 :     {
      51             219 :         return object_pool_.acquire(*this);
      52                 :     }
      53                 : 
      54             217 :     void destroy(io_object::implementation* p) override
      55                 :     {
      56             217 :         auto& impl = static_cast<posix_random_access_file&>(*p);
      57             217 :         impl.cancel();
      58             217 :         impl.close_file();
      59             217 :         release(&impl);
      60             217 :     }
      61                 : 
      62             413 :     void close(io_object::handle& h) override
      63                 :     {
      64             413 :         if (h.get())
      65                 :         {
      66             413 :             auto& impl = static_cast<posix_random_access_file&>(*h.get());
      67             413 :             impl.cancel();
      68             413 :             impl.close_file();
      69                 :         }
      70             413 :     }
      71                 : 
      72             203 :     std::error_code open_file(
      73                 :         random_access_file::implementation& impl,
      74                 :         std::filesystem::path const& path,
      75                 :         file_base::flags mode) override
      76                 :     {
      77                 :         // Unavailable in the unsafe tier: the file thread pool completes
      78                 :         // cross-thread, which the lockless scheduler cannot accept.
      79             203 :         if (sched_->scheduler_locking_disabled())
      80 MIS           0 :             return std::make_error_code(std::errc::operation_not_supported);
      81 HIT         203 :         return static_cast<posix_random_access_file&>(impl).open_file(
      82             203 :             path, mode);
      83                 :     }
      84                 : 
      85             164 :     void shutdown() override
      86                 :     {
      87                 :         // See uring_socket_service_base::shutdown(): snapshot under an
      88                 :         // acquired reference, cancel+close without the pool lock held;
      89                 :         // shutdown() sets shutting-down and takes the snapshot in one
      90                 :         // critical section, so a close that drops the last ref deletes
      91                 :         // rather than recycles.
      92             164 :         std::vector<posix_random_access_file*> live;
      93             164 :         object_pool_.shutdown(
      94             164 :             [&](posix_random_access_file* f)
      95                 :             {
      96               4 :                 acquire(f);
      97               4 :                 live.push_back(f);
      98               4 :             });
      99             168 :         for (auto* f : live)
     100                 :         {
     101               4 :             f->cancel();
     102               4 :             f->close_file();
     103               4 :             release(f);
     104                 :         }
     105             164 :     }
     106                 : 
     107             490 :     void post(scheduler_op* op)
     108                 :     {
     109             490 :         sched_->post(op);
     110             490 :     }
     111                 : 
     112                 :     void work_started() noexcept
     113                 :     {
     114                 :         sched_->work_started();
     115                 :     }
     116                 : 
     117                 :     void work_finished() noexcept
     118                 :     {
     119                 :         sched_->work_finished();
     120                 :     }
     121                 : 
     122                 :     /** Return the thread pool that runs this service's file work.
     123                 : 
     124                 :         The pool's service is created on first use, so this can fail
     125                 :         where a plain accessor could not. Its workers start later, on
     126                 :         the first post, and a thread the system refuses there is
     127                 :         reported by that post rather than thrown here.
     128                 : 
     129                 :         @throws std::bad_alloc If the service cannot be allocated.
     130                 : 
     131                 :         @return The context's shared blocking-I/O pool.
     132                 : 
     133                 :         @see thread_pool_ref::get
     134                 :     */
     135             494 :     thread_pool& pool()
     136                 :     {
     137             494 :         return pool_.get();
     138                 :     }
     139                 : 
     140                 : private:
     141                 :     scheduler* sched_;
     142                 :     thread_pool_ref pool_;
     143                 :     object_pool<posix_random_access_file> object_pool_;
     144                 : };
     145                 : 
     146                 : // ---------------------------------------------------------------------------
     147                 : // posix_random_access_file inline implementations (require complete service)
     148                 : // ---------------------------------------------------------------------------
     149                 : 
     150                 : inline std::coroutine_handle<>
     151             379 : posix_random_access_file::read_some_at(
     152                 :     std::uint64_t offset,
     153                 :     std::coroutine_handle<> h,
     154                 :     capy::executor_ref ex,
     155                 :     buffer_param param,
     156                 :     std::stop_token token,
     157                 :     std::error_code* ec,
     158                 :     std::size_t* bytes_out)
     159                 : {
     160                 :     // Closed-object contract outranks the zero-length no-op.
     161             379 :     if (fd_ < 0)
     162                 :     {
     163               4 :         *ec        = make_error_code(std::errc::bad_file_descriptor);
     164               4 :         *bytes_out = 0;
     165               4 :         return h;
     166                 :     }
     167                 : 
     168             375 :     capy::mutable_buffer bufs[max_buffers];
     169             375 :     auto count = param.copy_to(bufs, max_buffers);
     170                 : 
     171             375 :     if (count == 0)
     172                 :     {
     173               2 :         *ec        = {};
     174               2 :         *bytes_out = 0;
     175               2 :         return h;
     176                 :     }
     177                 : 
     178             373 :     auto* op    = acquire_op();
     179             373 :     op->is_read = true;
     180             373 :     op->offset  = offset;
     181                 : 
     182             373 :     op->iovec_count = static_cast<int>(count);
     183             746 :     for (int i = 0; i < op->iovec_count; ++i)
     184                 :     {
     185             373 :         op->iovecs[i].iov_base = bufs[i].data();
     186             373 :         op->iovecs[i].iov_len  = bufs[i].size();
     187                 :     }
     188                 : 
     189             373 :     op->h           = h;
     190             373 :     op->ex          = ex;
     191             373 :     op->ec_out      = ec;
     192             373 :     op->bytes_out   = bytes_out;
     193             373 :     op->file_       = this;
     194             373 :     op->object_ref_ = detail::object_ref(this);
     195             373 :     op->start(token);
     196                 : 
     197             373 :     op->ex.on_work_started();
     198                 : 
     199                 :     {
     200             373 :         std::lock_guard<std::mutex> lock(ops_mutex_);
     201             373 :         outstanding_ops_.push_back(op);
     202             373 :     }
     203                 : 
     204             373 :     static_cast<pool_work_item*>(op)->func_ = &raf_op::do_work;
     205             373 :     if (auto pec = svc_.pool().post(static_cast<pool_work_item*>(op)))
     206                 :     {
     207                 :         // The pool is shutting down, or the system refused it a thread.
     208                 :         // Nothing of this read went cross-thread, so it answers here
     209                 :         // like the closed-descriptor and zero-length exits above rather
     210                 :         // than through a completion the scheduler has to carry back.
     211                 :         // destroy() is the discard the op never reaching the queue
     212                 :         // needs: it unlinks, unwinds the work count and recycles.
     213               2 :         op->destroy();
     214               2 :         *ec        = pec;
     215               2 :         *bytes_out = 0;
     216               2 :         return h;
     217                 :     }
     218             371 :     return std::noop_coroutine();
     219                 : }
     220                 : 
     221                 : inline std::coroutine_handle<>
     222             125 : posix_random_access_file::write_some_at(
     223                 :     std::uint64_t offset,
     224                 :     std::coroutine_handle<> h,
     225                 :     capy::executor_ref ex,
     226                 :     buffer_param param,
     227                 :     std::stop_token token,
     228                 :     std::error_code* ec,
     229                 :     std::size_t* bytes_out)
     230                 : {
     231                 :     // Closed-object contract outranks the zero-length no-op.
     232             125 :     if (fd_ < 0)
     233                 :     {
     234               2 :         *ec        = make_error_code(std::errc::bad_file_descriptor);
     235               2 :         *bytes_out = 0;
     236               2 :         return h;
     237                 :     }
     238                 : 
     239             123 :     capy::mutable_buffer bufs[max_buffers];
     240             123 :     auto count = param.copy_to(bufs, max_buffers);
     241                 : 
     242             123 :     if (count == 0)
     243                 :     {
     244               2 :         *ec        = {};
     245               2 :         *bytes_out = 0;
     246               2 :         return h;
     247                 :     }
     248                 : 
     249             121 :     auto* op    = acquire_op();
     250             121 :     op->is_read = false;
     251             121 :     op->offset  = offset;
     252                 : 
     253             121 :     op->iovec_count = static_cast<int>(count);
     254             242 :     for (int i = 0; i < op->iovec_count; ++i)
     255                 :     {
     256             121 :         op->iovecs[i].iov_base = bufs[i].data();
     257             121 :         op->iovecs[i].iov_len  = bufs[i].size();
     258                 :     }
     259                 : 
     260             121 :     op->h           = h;
     261             121 :     op->ex          = ex;
     262             121 :     op->ec_out      = ec;
     263             121 :     op->bytes_out   = bytes_out;
     264             121 :     op->file_       = this;
     265             121 :     op->object_ref_ = detail::object_ref(this);
     266             121 :     op->start(token);
     267                 : 
     268             121 :     op->ex.on_work_started();
     269                 : 
     270                 :     {
     271             121 :         std::lock_guard<std::mutex> lock(ops_mutex_);
     272             121 :         outstanding_ops_.push_back(op);
     273             121 :     }
     274                 : 
     275             121 :     static_cast<pool_work_item*>(op)->func_ = &raf_op::do_work;
     276             121 :     if (auto pec = svc_.pool().post(static_cast<pool_work_item*>(op)))
     277                 :     {
     278                 :         // The pool is shutting down, or the system refused it a thread.
     279                 :         // Nothing of this write went cross-thread, so it answers here
     280                 :         // like the closed-descriptor and zero-length exits above rather
     281                 :         // than through a completion the scheduler has to carry back.
     282                 :         // destroy() is the discard the op never reaching the queue
     283                 :         // needs: it unlinks, unwinds the work count and recycles.
     284               2 :         op->destroy();
     285               2 :         *ec        = pec;
     286               2 :         *bytes_out = 0;
     287               2 :         return h;
     288                 :     }
     289             119 :     return std::noop_coroutine();
     290                 : }
     291                 : 
     292                 : // -- raf_op thread-pool work function --
     293                 : 
     294                 : inline void
     295             490 : posix_random_access_file::raf_op::do_work(pool_work_item* w) noexcept
     296                 : {
     297             490 :     auto* op   = static_cast<raf_op*>(w);
     298             490 :     auto* self = op->file_;
     299                 : 
     300             490 :     if (op->cancelled.load(std::memory_order_acquire))
     301                 :     {
     302              61 :         op->errn              = ECANCELED;
     303              61 :         op->bytes_transferred = 0;
     304                 :     }
     305             429 :     else if (
     306             858 :         op->offset >
     307             429 :         static_cast<std::uint64_t>(std::numeric_limits<off_t>::max()))
     308                 :     {
     309               2 :         op->errn              = EOVERFLOW;
     310               2 :         op->bytes_transferred = 0;
     311                 :     }
     312                 :     else
     313                 :     {
     314                 :         ssize_t n;
     315             427 :         if (op->is_read)
     316                 :         {
     317                 :             do
     318                 :             {
     319             634 :                 n = ::preadv(
     320             317 :                     self->fd_, op->iovecs, op->iovec_count,
     321             317 :                     static_cast<off_t>(op->offset));
     322                 :             }
     323             317 :             while (n < 0 && errno == EINTR);
     324                 :         }
     325                 :         else
     326                 :         {
     327                 :             do
     328                 :             {
     329             220 :                 n = ::pwritev(
     330             110 :                     self->fd_, op->iovecs, op->iovec_count,
     331             110 :                     static_cast<off_t>(op->offset));
     332                 :             }
     333             110 :             while (n < 0 && errno == EINTR);
     334                 :         }
     335                 : 
     336             427 :         if (n >= 0)
     337                 :         {
     338             413 :             op->errn              = 0;
     339             413 :             op->bytes_transferred = static_cast<std::size_t>(n);
     340                 :         }
     341                 :         else
     342                 :         {
     343              14 :             op->errn              = errno;
     344              14 :             op->bytes_transferred = 0;
     345                 :         }
     346                 :     }
     347                 : 
     348             490 :     self->svc_.post(static_cast<scheduler_op*>(op));
     349             490 : }
     350                 : 
     351                 : inline void
     352             217 : posix_random_access_file::retire() noexcept
     353                 : {
     354             217 :     svc_.object_pool_.recycle(this);
     355             217 : }
     356                 : 
     357                 : } // namespace boost::corosio::detail
     358                 : 
     359                 : #endif // BOOST_COROSIO_POSIX
     360                 : 
     361                 : #endif // BOOST_COROSIO_NATIVE_DETAIL_POSIX_POSIX_RANDOM_ACCESS_FILE_SERVICE_HPP
        

Generated by: LCOV version 2.3