#include <so_5/all.hpp>
using namespace std::chrono_literals;
class logger
{
std::mutex m_lock;
public:
template< typename... Args >
void info( Args && ...args )
{
std::lock_guard< std::mutex > lock{ m_lock };
std::cout << "*** ";
(std::cout << ... << args);
std::cout << std::endl;
}
template< typename... Args >
void err( Args && ...args )
{
std::lock_guard< std::mutex > lock{ m_lock };
std::cout << "### ";
(std::cout << ... << args);
std::cout << std::endl;
}
};
logger g_log;
class msg_do_something final : public so_5::message_t
{
private:
bool m_processed{ false };
public:
std::string m_id;
msg_do_something( std::string id ) : m_id{ std::move(id) } {}
~msg_do_something() noexcept override
{
if( !m_processed )
{
g_log.err( "[", m_id, "] discarded without processing" );
}
}
void processed() noexcept
{
m_processed = true;
}
};
class processor final : public so_5::agent_t
{
const so_5::mbox_t m_incoming_mbox;
const std::string m_name;
public:
processor(
context_t ctx,
so_5::mbox_t incoming_mbox,
std::string name )
:
so_5::agent_t{ std::move(ctx) }
, m_incoming_mbox{ std::move(incoming_mbox) }
, m_name{ std::move(name) }
{}
void so_define_agent() override
{
so_subscribe( m_incoming_mbox )
.event( &processor::evt_do_something );
}
private:
void evt_do_something( mutable_mhood_t<msg_do_something> cmd )
{
g_log.info( m_name, " [", cmd->m_id, "] processing started" );
std::this_thread::sleep_for( 25ms );
g_log.info( m_name, " [", cmd->m_id, "] processing finished" );
cmd->processed();
}
};
class generator final : public so_5::agent_t
{
struct msg_generate_next final : public so_5::signal_t {};
const std::string m_name;
const so_5::mbox_t m_dest_mbox;
const std::chrono::milliseconds m_initial_delay;
unsigned int m_ordinal{};
public:
generator(
context_t ctx,
std::string name,
so_5::mbox_t dest_mbox,
std::chrono::milliseconds initial_delay )
:
so_5::agent_t{ std::move(ctx) }
, m_name{ std::move(name) }
, m_dest_mbox{ std::move(dest_mbox) }
, m_initial_delay{ initial_delay }
{}
void so_define_agent() override
{
so_subscribe_self()
.event( &generator::evt_generate_next );
}
void so_evt_start() override
{
so_5::send_delayed< msg_generate_next >( *this, m_initial_delay );
}
private:
void evt_generate_next( mhood_t<msg_generate_next> )
{
++m_ordinal;
auto id = m_name + "-" + std::to_string( m_ordinal );
g_log.info( m_name, " sending [", id, "]" );
so_5::send< so_5::mutable_msg<msg_do_something> >(
m_dest_mbox,
std::move(id) );
so_5::send_delayed< msg_generate_next >( *this, 15ms );
}
};
int main()
{
so_5::launch( []( so_5::environment_t & env ) {
env.introduce_coop( []( so_5::coop_t & coop ) {
coop.environment() );
constexpr unsigned int processors_count = 4u;
auto thread_pool_binder =
so_5::disp::thread_pool::make_dispatcher(
coop.environment(),
processors_count )
.binder( []( auto & bind_params ) {
bind_params.fifo( so_5::disp::thread_pool::
fifo_t::individual );
} );
for( unsigned int i = 0; i != processors_count; ++i )
{
coop.make_agent_with_binder< processor >(
thread_pool_binder,
rr_mbox,
"worker-" + std::to_string( i + 1u ) );
}
auto dest_mbox = so_5::extra::mboxes::inflight_limit::
make_mbox< so_5::mutable_msg<msg_do_something> >(
rr_mbox,
processors_count );
constexpr std::size_t generators_count = 4u;
std::string names[ generators_count ] =
{ "alice", "bob", "eve", "kate" };
std::chrono::milliseconds initial_delays[ generators_count ] =
{ 7ms, 0ms, 17ms, 23ms };
for( std::size_t i{}; i != generators_count; ++i )
coop.make_agent< generator >(
names[ i ],
dest_mbox,
initial_delays[ i ] );
} );
std::this_thread::sleep_for( 95ms );
env.stop();
} );
}
Implementation of proxy mbox with inflight limit.
Ranges for error codes of each submodules.
Implementation of round-robin mbox.