Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 4316c3ce95 |
@@ -21,6 +21,7 @@
|
||||
#include <iostream>
|
||||
#include <algorithm>
|
||||
#include <mutex>
|
||||
#include <chrono>
|
||||
#include <boost/algorithm/string/predicate.hpp>
|
||||
#include <boost/range/adaptor/map.hpp>
|
||||
#include <boost/multi_index_container.hpp>
|
||||
@@ -434,8 +435,12 @@ class producer_plugin_impl : public std::enable_shared_from_this<producer_plugin
|
||||
fc::microseconds _ro_max_trx_time_us{ 0 }; // calculated during option initialization
|
||||
ro_trx_queue_t _ro_trx_queue;
|
||||
std::atomic<uint32_t> _ro_num_active_trx_exec_tasks{ 0 };
|
||||
std::vector<std::future<bool>> _ro_trx_exec_tasks_fut;
|
||||
std::condition_variable _switch_window_cond;
|
||||
std::mutex _switch_window_mtx;
|
||||
uint32_t _ro_idle_threads{ 0 };
|
||||
std::condition_variable _ro_more_work_cond;
|
||||
|
||||
void start_read_only_thread();
|
||||
void start_write_window();
|
||||
void switch_to_write_window();
|
||||
void switch_to_read_window();
|
||||
@@ -645,6 +650,9 @@ class producer_plugin_impl : public std::enable_shared_from_this<producer_plugin
|
||||
// executed in read window
|
||||
std::lock_guard<std::mutex> g( _ro_trx_queue.mtx );
|
||||
_ro_trx_queue.queue.push_back({trx, std::move(next)});
|
||||
if ( _ro_idle_threads > 0 ) {
|
||||
_ro_more_work_cond.notify_all();
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -1214,8 +1222,8 @@ void producer_plugin::plugin_startup()
|
||||
fc_elog( _log, "Exception in read-only transaction thread pool, exiting: ${e}", ("e", e.to_detail_string()) );
|
||||
app().quit();
|
||||
},
|
||||
[&]() {
|
||||
chain.init_thread_local_data();
|
||||
[this]() {
|
||||
this->my->start_read_only_thread();
|
||||
});
|
||||
|
||||
my->start_write_window();
|
||||
@@ -2772,7 +2780,7 @@ void producer_plugin_impl::switch_to_write_window() {
|
||||
return;
|
||||
}
|
||||
|
||||
EOS_ASSERT(_ro_num_active_trx_exec_tasks.load() == 0 && _ro_trx_exec_tasks_fut.empty(), producer_exception, "no read-only tasks should be running before switching to write window");
|
||||
EOS_ASSERT(_ro_num_active_trx_exec_tasks.load() == 0, producer_exception, "no read-only tasks should be running before switching to write window");
|
||||
_ro_read_window_timer.cancel();
|
||||
_ro_write_window_timer.cancel();
|
||||
|
||||
@@ -2786,6 +2794,10 @@ void producer_plugin_impl::start_write_window() {
|
||||
app().executor().set_to_write_window();
|
||||
chain.unset_db_read_only_mode();
|
||||
_idle_trx_time = fc::time_point::now();
|
||||
{
|
||||
std::lock_guard<std::mutex> g(_switch_window_mtx);
|
||||
_switch_window_cond.notify_all();
|
||||
}
|
||||
|
||||
auto expire_time = boost::posix_time::microseconds(_ro_write_window_time_us.count());
|
||||
_ro_write_window_timer.expires_from_now( expire_time );
|
||||
@@ -2803,7 +2815,7 @@ void producer_plugin_impl::start_write_window() {
|
||||
// Called from app thread
|
||||
void producer_plugin_impl::switch_to_read_window() {
|
||||
EOS_ASSERT(app().executor().is_write_window(), producer_exception, "expected to be in write window");
|
||||
EOS_ASSERT(_ro_num_active_trx_exec_tasks.load() == 0 && _ro_trx_exec_tasks_fut.empty(), producer_exception, "_ro_trx_exec_tasks_fut expected to be empty" );
|
||||
EOS_ASSERT(_ro_num_active_trx_exec_tasks.load() == 0, producer_exception, "_ro_num_active_trx_exec_tasks expected to be 0" );
|
||||
|
||||
_ro_write_window_timer.cancel();
|
||||
_ro_read_window_timer.cancel();
|
||||
@@ -2819,33 +2831,37 @@ void producer_plugin_impl::switch_to_read_window() {
|
||||
app().executor().set_to_read_window();
|
||||
chain_plug->chain().set_db_read_only_mode();
|
||||
_received_block = false;
|
||||
|
||||
// start a read-only transaction execution task in each thread in the thread pool
|
||||
_ro_num_active_trx_exec_tasks = _ro_thread_pool_size;
|
||||
for (auto i = 0; i < _ro_thread_pool_size; ++i ) {
|
||||
_ro_trx_exec_tasks_fut.emplace_back( post_async_task( _ro_thread_pool.get_executor(), [self = this] () {
|
||||
return self->read_only_trx_execution_task();
|
||||
}) );
|
||||
{
|
||||
std::lock_guard<std::mutex> g(_switch_window_mtx);
|
||||
_switch_window_cond.notify_all();
|
||||
}
|
||||
// read-only threads decide when to switch back to write window
|
||||
}
|
||||
|
||||
auto expire_time = boost::posix_time::microseconds(_ro_read_window_time_us.count());
|
||||
_ro_read_window_timer.expires_from_now( expire_time );
|
||||
_ro_read_window_timer.async_wait( app().executor().wrap( // stay on app thread
|
||||
priority::high,
|
||||
exec_queue::read_only_trx_safe,
|
||||
[weak_this = weak_from_this()]( const boost::system::error_code& ec ) {
|
||||
auto self = weak_this.lock();
|
||||
if( self && ec != boost::asio::error::operation_aborted ) {
|
||||
// use future to make sure all read-only tasks finished before switching to write window
|
||||
for ( auto& task: self->_ro_trx_exec_tasks_fut ) {
|
||||
task.get();
|
||||
}
|
||||
self->_ro_trx_exec_tasks_fut.clear();
|
||||
self->switch_to_write_window();
|
||||
} else if ( self ) {
|
||||
self->_ro_trx_exec_tasks_fut.clear();
|
||||
}
|
||||
}));
|
||||
// Called from a read only trx thread. Run in parallel with main and other read only trx threads
|
||||
void producer_plugin_impl::start_read_only_thread() {
|
||||
chain::controller& chain = chain_plug->chain();
|
||||
chain.init_thread_local_data();
|
||||
|
||||
while ( 1 ) {
|
||||
{
|
||||
// wait until we are in read window
|
||||
std::unique_lock<std::mutex> lck(_switch_window_mtx);
|
||||
_switch_window_cond.wait(lck, []{ return app().executor().is_read_window(); });
|
||||
}
|
||||
|
||||
// executue trxs until deadline or no more trxs
|
||||
read_only_trx_execution_task();
|
||||
|
||||
if ( --_ro_num_active_trx_exec_tasks == 0 ) {
|
||||
// last active thread switches the window
|
||||
switch_to_write_window();
|
||||
} else {
|
||||
// other threads wait
|
||||
std::unique_lock<std::mutex> lck(_switch_window_mtx);
|
||||
_switch_window_cond.wait(lck, []{ return app().executor().is_write_window(); });
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Called from a read only trx thread. Run in parallel with app and other read only trx threads
|
||||
@@ -2860,7 +2876,27 @@ bool producer_plugin_impl::read_only_trx_execution_task() {
|
||||
while ( fc::time_point::now() < read_window_deadline && !_received_block ) {
|
||||
std::unique_lock<std::mutex> lck( _ro_trx_queue.mtx );
|
||||
if ( _ro_trx_queue.queue.empty() ) {
|
||||
break;
|
||||
// double check so read_window_deadline - now is valid
|
||||
auto now = fc::time_point::now();
|
||||
if ( now >= read_window_deadline || _received_block ) {
|
||||
break;
|
||||
}
|
||||
auto wait_time = std::min(read_window_deadline - now, _ro_max_trx_time_us).count();
|
||||
|
||||
++_ro_idle_threads;
|
||||
_ro_more_work_cond.wait_for(lck, std::chrono::microseconds(wait_time));
|
||||
--_ro_idle_threads;
|
||||
if ( fc::time_point::now() >= read_window_deadline || !_received_block ) {
|
||||
// honor read-window deadline and _received_block
|
||||
break;
|
||||
}
|
||||
if ( !_ro_trx_queue.queue.empty() || _ro_idle_threads > 0 ) {
|
||||
// may have more work
|
||||
continue;
|
||||
} else {
|
||||
// no need to wait
|
||||
break;
|
||||
}
|
||||
}
|
||||
auto trx = _ro_trx_queue.queue.front();
|
||||
_ro_trx_queue.queue.pop_front();
|
||||
@@ -2870,20 +2906,11 @@ bool producer_plugin_impl::read_only_trx_execution_task() {
|
||||
if ( retry ) {
|
||||
lck.lock();
|
||||
_ro_trx_queue.queue.push_front(trx);
|
||||
// Do not schedule new execution
|
||||
// This happens only when a trx exhausted its time. Do not schedule new executions
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
// If all tasks are finished, do not wait until end of read window; switch to write window now.
|
||||
if ( --_ro_num_active_trx_exec_tasks == 0 ) {
|
||||
// Do switching on app thread to serialize
|
||||
app().executor().post( priority::high, exec_queue::read_only_trx_safe, [self=this]() {
|
||||
self->_ro_trx_exec_tasks_fut.clear();
|
||||
self->switch_to_write_window();
|
||||
} );
|
||||
}
|
||||
|
||||
return true;
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user