Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
28 changes: 25 additions & 3 deletions include/boost/corosio/detail/scheduler.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -13,12 +13,14 @@
#define BOOST_COROSIO_DETAIL_SCHEDULER_HPP

#include <boost/corosio/detail/config.hpp>
#include <boost/corosio/detail/except.hpp>

#include <system_error>
#include <boost/capy/continuation.hpp>
#include <coroutine>
#include <boost/capy/ex/execution_context.hpp>

#include <coroutine>
#include <cstddef>
#include <system_error>

namespace boost::corosio::detail {

Expand All @@ -30,11 +32,18 @@ class scheduler_op;
this to implement the reactor/proactor event loop. The
@ref io_context delegates all scheduling operations here.

The scheduler is a registry service keyed under this abstract
type, so services created on first use can locate it without
naming a concrete backend.

@see io_context
*/
struct BOOST_COROSIO_DECL scheduler
: capy::execution_context::service
{
virtual ~scheduler() = default;
using key_type = scheduler;

~scheduler() override = default;

/// Post a coroutine handle for deferred execution.
virtual void post(std::coroutine_handle<>) const = 0;
Expand Down Expand Up @@ -126,6 +135,19 @@ struct BOOST_COROSIO_DECL scheduler
virtual void configure_threading(threading_config) noexcept = 0;
};

/** Return the scheduler registered with the context.

@throws std::logic_error If the context has no backend installed.
*/
inline scheduler&
get_scheduler(capy::execution_context& ctx)
{
auto* sched = ctx.find_service<scheduler>();
if (!sched)
throw_logic_error("no scheduler installed");
return *sched;
}

} // namespace boost::corosio::detail

#endif
4 changes: 3 additions & 1 deletion include/boost/corosio/native/detail/endpoint_convert.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -298,7 +298,9 @@ to_sockaddr(local_endpoint const& ep, sockaddr_storage& storage) noexcept
un_sa_t sa{};
sa.sun_family = AF_UNIX;
auto path = ep.path();
auto copy_len = (std::min)(path.size(), sizeof(sa.sun_path));
auto copy_len = (std::min)(
path.size(),
(std::min)(local_endpoint::max_path_length, sizeof(sa.sun_path)));
if (copy_len > 0)
std::memcpy(sa.sun_path, path.data(), copy_len);
std::memcpy(&storage, &sa, sizeof(sa));
Expand Down
9 changes: 0 additions & 9 deletions include/boost/corosio/native/detail/epoll/epoll_scheduler.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -24,10 +24,6 @@
#include <boost/corosio/native/detail/epoll/epoll_traits.hpp>
#include <boost/corosio/detail/timer_service.hpp>
#include <boost/corosio/native/detail/make_err.hpp>
#include <boost/corosio/native/detail/posix/posix_resolver_service.hpp>
#include <boost/corosio/native/detail/posix/posix_signal_service.hpp>
#include <boost/corosio/native/detail/posix/posix_stream_file_service.hpp>
#include <boost/corosio/native/detail/posix/posix_random_access_file_service.hpp>

#include <boost/corosio/detail/except.hpp>

Expand Down Expand Up @@ -215,11 +211,6 @@ inline epoll_scheduler::epoll_scheduler(capy::execution_context& ctx, int)
self->interrupt_reactor();
}));

get_resolver_service(ctx, *this);
get_signal_service(ctx, *this);
get_stream_file_service(ctx, *this);
get_random_access_file_service(ctx, *this);

completed_ops_.push(&task_op_);
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -43,8 +43,8 @@ class BOOST_COROSIO_DECL win_local_stream_acceptor_service final
: public local_stream_acceptor_service
{
public:
win_local_stream_acceptor_service(
capy::execution_context& ctx, win_local_stream_service& svc);
explicit win_local_stream_acceptor_service(
capy::execution_context& ctx);

io_object::implementation* construct() override;

Expand Down Expand Up @@ -577,8 +577,8 @@ win_local_stream_acceptor::get_internal() const noexcept
// ============================================================

inline win_local_stream_acceptor_service::win_local_stream_acceptor_service(
capy::execution_context& /*ctx*/, win_local_stream_service& svc)
: svc_(svc)
capy::execution_context& ctx)
: svc_(ctx.use_service<win_local_stream_service>())
{
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@
#include <boost/corosio/native/detail/iocp/win_dissociate.hpp>
#include <boost/corosio/native/detail/iocp/win_local_stream_acceptor.hpp>
#include <boost/corosio/native/detail/iocp/win_local_stream_socket.hpp>
#include <boost/corosio/native/detail/iocp/win_tcp_service.hpp>
#include <boost/corosio/native/detail/iocp/win_tcp_acceptor_service.hpp>
#include <boost/corosio/native/detail/iocp/win_scheduler.hpp>
#include <boost/corosio/native/detail/iocp/win_completion_key.hpp>
#include <boost/corosio/native/detail/iocp/win_mutex.hpp>
Expand Down Expand Up @@ -55,8 +55,7 @@ class BOOST_COROSIO_DECL win_local_stream_service final

void close(io_object::handle& h) override;

explicit win_local_stream_service(
capy::execution_context& ctx, win_tcp_service& tcp_svc);
explicit win_local_stream_service(capy::execution_context& ctx);

~win_local_stream_service();

Expand Down Expand Up @@ -877,8 +876,8 @@ win_local_stream_socket::get_internal() const noexcept
// ============================================================

inline win_local_stream_service::win_local_stream_service(
capy::execution_context& ctx, win_tcp_service& tcp_svc)
: tcp_svc_(tcp_svc)
capy::execution_context& ctx)
: tcp_svc_(ctx.use_service<win_tcp_service>())
, sched_(ctx.use_service<win_scheduler>())
, iocp_(sched_.native_handle())
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
#if BOOST_COROSIO_HAS_IOCP

#include <boost/corosio/native/detail/iocp/win_resolver.hpp>
#include <boost/corosio/detail/scheduler.hpp>
#include <boost/corosio/detail/thread_pool.hpp>

#include <unordered_map>
Expand Down Expand Up @@ -60,9 +61,8 @@ class BOOST_COROSIO_DECL win_resolver_service final
/** Construct the resolver service.

@param ctx Reference to the owning execution_context.
@param sched Reference to the scheduler for posting completions.
*/
win_resolver_service(capy::execution_context& ctx, scheduler& sched);
explicit win_resolver_service(capy::execution_context& ctx);

/** Destroy the resolver service. */
~win_resolver_service();
Expand Down Expand Up @@ -580,8 +580,8 @@ win_resolver::do_reverse_resolve_work(pool_work_item* w) noexcept
// win_resolver_service

inline win_resolver_service::win_resolver_service(
capy::execution_context& ctx, scheduler& sched)
: sched_(sched)
capy::execution_context& ctx)
: sched_(get_scheduler(ctx))
, pool_(ctx)
{
}
Expand Down
4 changes: 0 additions & 4 deletions include/boost/corosio/native/detail/iocp/win_scheduler.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,6 @@
#include <boost/corosio/native/detail/iocp/win_overlapped_op.hpp>
#include <boost/corosio/native/detail/iocp/win_timers.hpp>
#include <boost/corosio/detail/timer_service.hpp>
#include <boost/corosio/native/detail/iocp/win_resolver_service.hpp>
#include <boost/corosio/native/detail/make_err.hpp>
#include <boost/corosio/detail/except.hpp>
#include <boost/corosio/detail/thread_local_ptr.hpp>
Expand All @@ -53,10 +52,8 @@ class win_wait_reactor;

class BOOST_COROSIO_DECL win_scheduler final
: public scheduler
, public capy::execution_context::service
{
public:
using key_type = scheduler;

win_scheduler(capy::execution_context& ctx, int concurrency_hint = -1);
~win_scheduler();
Expand Down Expand Up @@ -726,7 +723,6 @@ inline win_scheduler::win_scheduler(
{
timers_ = make_win_timers(iocp_, &dispatch_required_);
set_timer_service(&get_timer_service(ctx, *this));
ctx.make_service<win_resolver_service>(*this);

// A scheduler whose wait reactor could not be built would
// answer every wait with a parked op, so it refuses to exist
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -47,8 +47,7 @@ class BOOST_COROSIO_DECL win_tcp_acceptor_service final
public:
using key_type = win_tcp_acceptor_service;

win_tcp_acceptor_service(
capy::execution_context& ctx, win_tcp_service& svc);
explicit win_tcp_acceptor_service(capy::execution_context& ctx);

io_object::implementation* construct() override;

Expand Down Expand Up @@ -1737,8 +1736,8 @@ win_tcp_acceptor::get_internal() const noexcept
// win_tcp_acceptor_service

inline win_tcp_acceptor_service::win_tcp_acceptor_service(
[[maybe_unused]] capy::execution_context& ctx, win_tcp_service& svc)
: svc_(svc)
capy::execution_context& ctx)
: svc_(ctx.use_service<win_tcp_service>())
{
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,10 +24,6 @@
#include <boost/corosio/native/detail/kqueue/kqueue_traits.hpp>
#include <boost/corosio/detail/timer_service.hpp>
#include <boost/corosio/native/detail/make_err.hpp>
#include <boost/corosio/native/detail/posix/posix_resolver_service.hpp>
#include <boost/corosio/native/detail/posix/posix_signal_service.hpp>
#include <boost/corosio/native/detail/posix/posix_stream_file_service.hpp>
#include <boost/corosio/native/detail/posix/posix_random_access_file_service.hpp>

#include <boost/corosio/detail/except.hpp>

Expand Down Expand Up @@ -201,11 +197,6 @@ inline kqueue_scheduler::kqueue_scheduler(capy::execution_context& ctx, int)
static_cast<kqueue_scheduler*>(p)->interrupt_reactor();
}));

get_resolver_service(ctx, *this);
get_signal_service(ctx, *this);
get_stream_file_service(ctx, *this);
get_random_access_file_service(ctx, *this);

completed_ops_.push(&task_op_);
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,9 +30,8 @@ class BOOST_COROSIO_DECL posix_random_access_file_service final
: public random_access_file_service
{
public:
posix_random_access_file_service(
capy::execution_context& ctx, scheduler& sched)
: sched_(&sched)
explicit posix_random_access_file_service(capy::execution_context& ctx)
: sched_(&get_scheduler(ctx))
, pool_(ctx)
{
}
Expand Down Expand Up @@ -152,13 +151,6 @@ class BOOST_COROSIO_DECL posix_random_access_file_service final
file_ptrs_;
};

/** Get or create the random-access file service for the given context. */
inline posix_random_access_file_service&
get_random_access_file_service(capy::execution_context& ctx, scheduler& sched)
{
return ctx.make_service<posix_random_access_file_service>(sched);
}

// ---------------------------------------------------------------------------
// posix_random_access_file inline implementations (require complete service)
// ---------------------------------------------------------------------------
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -35,8 +35,8 @@ class BOOST_COROSIO_DECL posix_resolver_service final
public:
using key_type = posix_resolver_service;

posix_resolver_service(capy::execution_context& ctx, scheduler& sched)
: sched_(&sched)
explicit posix_resolver_service(capy::execution_context& ctx)
: sched_(&get_scheduler(ctx))
, pool_(ctx)
{
}
Expand Down Expand Up @@ -95,18 +95,6 @@ class BOOST_COROSIO_DECL posix_resolver_service final
resolver_ptrs_;
};

/** Get or create the resolver service for the given context.

This function is called by the concrete scheduler during initialization
to create the resolver service with a reference to itself.

@param ctx Reference to the owning execution_context.
@param sched Reference to the scheduler for posting completions.
@return Reference to the resolver service.
*/
posix_resolver_service&
get_resolver_service(capy::execution_context& ctx, scheduler& sched);

// ---------------------------------------------------------------------------
// Inline implementation
// ---------------------------------------------------------------------------
Expand Down Expand Up @@ -608,14 +596,6 @@ posix_resolver_service::post(scheduler_op* op)
sched_->post(op);
}

// Free function to get/create the resolver service

inline posix_resolver_service&
get_resolver_service(capy::execution_context& ctx, scheduler& sched)
{
return ctx.make_service<posix_resolver_service>(sched);
}

} // namespace boost::corosio::detail

#endif // BOOST_COROSIO_POSIX
Expand Down
30 changes: 5 additions & 25 deletions include/boost/corosio/native/detail/posix/posix_signal_service.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -37,8 +37,8 @@

Concrete signal service implementation for POSIX backends. Manages signal
registrations via sigaction() and dispatches completions through the
scheduler. One instance per execution_context, created by
get_signal_service().
scheduler. One instance per execution_context, created on first use
by the public signal_set.

See the block comment further down for the full architecture overview.
*/
Expand Down Expand Up @@ -157,7 +157,7 @@ class BOOST_COROSIO_DECL posix_signal_service final
public:
using key_type = posix_signal_service;

posix_signal_service(capy::execution_context& ctx, scheduler& sched);
explicit posix_signal_service(capy::execution_context& ctx);
~posix_signal_service() override;

posix_signal_service(posix_signal_service const&) = delete;
Expand Down Expand Up @@ -246,18 +246,6 @@ class BOOST_COROSIO_DECL posix_signal_service final
posix_signal_service* prev_ = nullptr;
};

/** Get or create the signal service for the given context.

This function is called by the concrete scheduler during initialization
to create the signal service with a reference to itself.

@param ctx Reference to the owning execution_context.
@param sched Reference to the scheduler for posting completions.
@return Reference to the signal service.
*/
posix_signal_service&
get_signal_service(capy::execution_context& ctx, scheduler& sched);

} // namespace detail

} // namespace boost::corosio
Expand Down Expand Up @@ -519,8 +507,8 @@ posix_signal::cancel() noexcept
// posix_signal_service implementation

inline posix_signal_service::posix_signal_service(
capy::execution_context&, scheduler& sched)
: sched_(&sched)
capy::execution_context& ctx)
: sched_(&get_scheduler(ctx))
{
for (int i = 0; i < max_signal_number; ++i)
{
Expand Down Expand Up @@ -1048,14 +1036,6 @@ posix_signal_service::remove_service(posix_signal_service* service)
}
}

// get_signal_service - factory function

inline posix_signal_service&
get_signal_service(capy::execution_context& ctx, scheduler& sched)
{
return ctx.make_service<posix_signal_service>(sched);
}

} // namespace detail
} // namespace boost::corosio

Expand Down
Loading
Loading