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 % 152 151 1
Test Date: 2026-10-08 17:58:26 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         170 :     explicit posix_random_access_file_service(capy::execution_context& ctx)
      34             510 :         : sched_(&get_scheduler(ctx))
      35             170 :         , pool_(ctx)
      36                 :     {
      37             170 :     }
      38                 : 
      39             340 :     ~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             225 :     io_object::implementation* construct() override
      47                 :     {
      48             225 :         auto ptr   = std::make_shared<posix_random_access_file>(*this);
      49             225 :         auto* impl = ptr.get();
      50                 : 
      51                 :         {
      52             225 :             std::lock_guard<std::mutex> lock(mutex_);
      53             225 :             file_list_.push_back(impl);
      54             225 :             file_ptrs_[impl] = std::move(ptr);
      55             225 :         }
      56                 : 
      57             225 :         return impl;
      58             225 :     }
      59                 : 
      60             223 :     void destroy(io_object::implementation* p) override
      61                 :     {
      62             223 :         auto& impl = static_cast<posix_random_access_file&>(*p);
      63             223 :         impl.cancel();
      64             223 :         impl.close_file();
      65             223 :         destroy_impl(impl);
      66             223 :     }
      67                 : 
      68             417 :     void close(io_object::handle& h) override
      69                 :     {
      70             417 :         if (h.get())
      71                 :         {
      72             417 :             auto& impl = static_cast<posix_random_access_file&>(*h.get());
      73             417 :             impl.cancel();
      74             417 :             impl.close_file();
      75                 :         }
      76             417 :     }
      77                 : 
      78             207 :     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             207 :         if (sched_->scheduler_locking_disabled())
      86 MIS           0 :             return std::make_error_code(std::errc::operation_not_supported);
      87 HIT         207 :         return static_cast<posix_random_access_file&>(impl).open_file(
      88             207 :             path, mode);
      89                 :     }
      90                 : 
      91             170 :     void shutdown() override
      92                 :     {
      93             170 :         std::lock_guard<std::mutex> lock(mutex_);
      94             172 :         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             170 :         file_ptrs_.clear();
     101             170 :     }
     102                 : 
     103             223 :     void destroy_impl(posix_random_access_file& impl)
     104                 :     {
     105             223 :         std::lock_guard<std::mutex> lock(mutex_);
     106             223 :         file_list_.remove(&impl);
     107             223 :         file_ptrs_.erase(&impl);
     108             223 :     }
     109                 : 
     110             432 :     void post(scheduler_op* op)
     111                 :     {
     112             432 :         sched_->post(op);
     113             432 :     }
     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             436 :     thread_pool& pool()
     139                 :     {
     140             436 :         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             349 : posix_random_access_file::read_some_at(
     160                 :     std::uint64_t offset,
     161                 :     capy::continuation& cont,
     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             349 :     if (fd_ < 0)
     170                 :     {
     171               4 :         *ec        = make_error_code(std::errc::bad_file_descriptor);
     172               4 :         *bytes_out = 0;
     173               4 :         return cont.h;
     174                 :     }
     175                 : 
     176             345 :     capy::mutable_buffer bufs[max_buffers];
     177             345 :     auto count = param.copy_to(bufs, max_buffers);
     178                 : 
     179             345 :     if (count == 0)
     180                 :     {
     181               2 :         *ec        = {};
     182               2 :         *bytes_out = 0;
     183               2 :         return cont.h;
     184                 :     }
     185                 : 
     186             343 :     auto* op    = new raf_op();
     187             343 :     op->is_read = true;
     188             343 :     op->offset  = offset;
     189             343 :     op->fd      = fd_;
     190                 : 
     191             343 :     op->iovec_count = static_cast<int>(count);
     192             686 :     for (int i = 0; i < op->iovec_count; ++i)
     193                 :     {
     194             343 :         op->iovecs[i].iov_base = bufs[i].data();
     195             343 :         op->iovecs[i].iov_len  = bufs[i].size();
     196                 :     }
     197                 : 
     198             343 :     op->h         = cont.h;
     199             343 :     op->awaiting  = &cont;
     200             343 :     op->ex        = ex;
     201             343 :     op->ec_out    = ec;
     202             343 :     op->bytes_out = bytes_out;
     203             343 :     op->file_     = this;
     204             343 :     op->impl_ptr  = this->shared_from_this();
     205             343 :     op->start(token);
     206                 : 
     207             343 :     op->ex.on_work_started();
     208                 : 
     209                 :     {
     210             343 :         std::lock_guard<std::mutex> lock(ops_mutex_);
     211             343 :         outstanding_ops_.push_back(op);
     212             343 :     }
     213                 : 
     214             343 :     static_cast<pool_work_item*>(op)->func_ = &raf_op::do_work;
     215             343 :     if (auto pec = svc_.pool().post(static_cast<pool_work_item*>(op)))
     216                 :     {
     217                 :         // The pool is shutting down, or the system refused it a thread.
     218                 :         // Nothing of this read went cross-thread, so it answers here
     219                 :         // like the closed-descriptor and zero-length exits above rather
     220                 :         // than through a completion the scheduler has to carry back.
     221                 :         // destroy() is the discard the op never reaching the queue
     222                 :         // needs: it unlinks, unwinds the work count and frees.
     223               2 :         op->destroy();
     224               2 :         *ec        = pec;
     225               2 :         *bytes_out = 0;
     226               2 :         return cont.h;
     227                 :     }
     228             341 :     return std::noop_coroutine();
     229                 : }
     230                 : 
     231                 : inline std::coroutine_handle<>
     232              97 : posix_random_access_file::write_some_at(
     233                 :     std::uint64_t offset,
     234                 :     capy::continuation& cont,
     235                 :     capy::executor_ref ex,
     236                 :     buffer_param param,
     237                 :     std::stop_token token,
     238                 :     std::error_code* ec,
     239                 :     std::size_t* bytes_out)
     240                 : {
     241                 :     // Closed-object contract outranks the zero-length no-op.
     242              97 :     if (fd_ < 0)
     243                 :     {
     244               2 :         *ec        = make_error_code(std::errc::bad_file_descriptor);
     245               2 :         *bytes_out = 0;
     246               2 :         return cont.h;
     247                 :     }
     248                 : 
     249              95 :     capy::mutable_buffer bufs[max_buffers];
     250              95 :     auto count = param.copy_to(bufs, max_buffers);
     251                 : 
     252              95 :     if (count == 0)
     253                 :     {
     254               2 :         *ec        = {};
     255               2 :         *bytes_out = 0;
     256               2 :         return cont.h;
     257                 :     }
     258                 : 
     259              93 :     auto* op    = new raf_op();
     260              93 :     op->is_read = false;
     261              93 :     op->offset  = offset;
     262              93 :     op->fd      = fd_;
     263                 : 
     264              93 :     op->iovec_count = static_cast<int>(count);
     265             186 :     for (int i = 0; i < op->iovec_count; ++i)
     266                 :     {
     267              93 :         op->iovecs[i].iov_base = bufs[i].data();
     268              93 :         op->iovecs[i].iov_len  = bufs[i].size();
     269                 :     }
     270                 : 
     271              93 :     op->h         = cont.h;
     272              93 :     op->awaiting  = &cont;
     273              93 :     op->ex        = ex;
     274              93 :     op->ec_out    = ec;
     275              93 :     op->bytes_out = bytes_out;
     276              93 :     op->file_     = this;
     277              93 :     op->impl_ptr  = this->shared_from_this();
     278              93 :     op->start(token);
     279                 : 
     280              93 :     op->ex.on_work_started();
     281                 : 
     282                 :     {
     283              93 :         std::lock_guard<std::mutex> lock(ops_mutex_);
     284              93 :         outstanding_ops_.push_back(op);
     285              93 :     }
     286                 : 
     287              93 :     static_cast<pool_work_item*>(op)->func_ = &raf_op::do_work;
     288              93 :     if (auto pec = svc_.pool().post(static_cast<pool_work_item*>(op)))
     289                 :     {
     290                 :         // The pool is shutting down, or the system refused it a thread.
     291                 :         // Nothing of this write went cross-thread, so it answers here
     292                 :         // like the closed-descriptor and zero-length exits above rather
     293                 :         // than through a completion the scheduler has to carry back.
     294                 :         // destroy() is the discard the op never reaching the queue
     295                 :         // needs: it unlinks, unwinds the work count and frees.
     296               2 :         op->destroy();
     297               2 :         *ec        = pec;
     298               2 :         *bytes_out = 0;
     299               2 :         return cont.h;
     300                 :     }
     301              91 :     return std::noop_coroutine();
     302                 : }
     303                 : 
     304                 : // -- raf_op thread-pool work function --
     305                 : 
     306                 : inline void
     307             432 : posix_random_access_file::raf_op::do_work(pool_work_item* w) noexcept
     308                 : {
     309             432 :     auto* op   = static_cast<raf_op*>(w);
     310             432 :     auto* self = op->file_;
     311                 : 
     312             432 :     if (op->cancelled.load(std::memory_order_acquire))
     313                 :     {
     314             102 :         op->errn              = ECANCELED;
     315             102 :         op->bytes_transferred = 0;
     316                 :     }
     317             330 :     else if (
     318             660 :         op->offset >
     319             330 :         static_cast<std::uint64_t>(std::numeric_limits<file_off_t>::max()))
     320                 :     {
     321               2 :         op->errn              = EOVERFLOW;
     322               2 :         op->bytes_transferred = 0;
     323                 :     }
     324                 :     else
     325                 :     {
     326                 :         ssize_t n;
     327             328 :         if (op->is_read)
     328                 :         {
     329                 :             do
     330                 :             {
     331             287 :                 n = file_preadv(
     332             287 :                     op->fd, op->iovecs, op->iovec_count,
     333             287 :                     static_cast<file_off_t>(op->offset));
     334                 :             }
     335             287 :             while (n < 0 && errno == EINTR);
     336                 :         }
     337                 :         else
     338                 :         {
     339                 :             do
     340                 :             {
     341              41 :                 n = file_pwritev(
     342              41 :                     op->fd, op->iovecs, op->iovec_count,
     343              41 :                     static_cast<file_off_t>(op->offset));
     344                 :             }
     345              41 :             while (n < 0 && errno == EINTR);
     346                 :         }
     347                 : 
     348             328 :         if (n >= 0)
     349                 :         {
     350             314 :             op->errn              = 0;
     351             314 :             op->bytes_transferred = static_cast<std::size_t>(n);
     352                 :         }
     353                 :         else
     354                 :         {
     355              14 :             op->errn              = errno;
     356              14 :             op->bytes_transferred = 0;
     357                 :         }
     358                 :     }
     359                 : 
     360             432 :     self->svc_.post(static_cast<scheduler_op*>(op));
     361             432 : }
     362                 : 
     363                 : } // namespace boost::corosio::detail
     364                 : 
     365                 : #endif // BOOST_COROSIO_POSIX
     366                 : 
     367                 : #endif // BOOST_COROSIO_NATIVE_DETAIL_POSIX_POSIX_RANDOM_ACCESS_FILE_SERVICE_HPP
        

Generated by: LCOV version 2.3