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 % 148 147 1
Test Date: 2026-09-28 20:06:38 Functions: 100.0 % 14 14

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

Generated by: LCOV version 2.3