diff --git a/include/atframe/modules/worker_context.h b/include/atframe/modules/worker_context.h index 123c3f3..1c4f1f2 100644 --- a/include/atframe/modules/worker_context.h +++ b/include/atframe/modules/worker_context.h @@ -12,15 +12,22 @@ #include #include -#include +#include LIBATAPP_MACRO_NAMESPACE_BEGIN struct UTIL_SYMBOL_VISIBLE worker_context { + // worker id 指示当前是第几个worker,0表示主线程,1表示第一个工作线程,依次类推。 + // worker id 可能被复用或转移工作线程,但同时每个 worker id 指向唯一一个线程 uint32_t worker_id = 0; - inline worker_context() noexcept : worker_id(0) {} - explicit inline worker_context(uint32_t id) noexcept : worker_id(id) {} + // worker_unique_id 指示当前worker的唯一标识,不会随着线程转移而变化 + // 注意: 在foreach接口中,如果对于stable的worker尚未分配完成,这个值可能传0 + uint64_t worker_unique_id = 0; + + inline worker_context() noexcept : worker_id(0), worker_unique_id(0) {} + explicit inline worker_context(uint32_t id, uint64_t unique_id) noexcept + : worker_id(id), worker_unique_id(unique_id) {} }; enum class worker_job_event_type : uint32_t { @@ -44,6 +51,11 @@ struct UTIL_SYMBOL_VISIBLE worker_meta { using worker_job_action_type = std::function; +using worker_event_callback_type = std::function; + +struct UTIL_SYMBOL_VISIBLE worker_event_callback_handle_data; +using worker_event_callback_handle_type = std::shared_ptr; + using worker_job_action_pointer = ::atfw::util::memory::strong_rc_ptr; struct UTIL_SYMBOL_VISIBLE worker_job_data { @@ -74,4 +86,3 @@ enum class worker_type : int32_t { }; LIBATAPP_MACRO_NAMESPACE_END - diff --git a/include/atframe/modules/worker_pool_module.h b/include/atframe/modules/worker_pool_module.h index 564b3b7..7f93a03 100644 --- a/include/atframe/modules/worker_pool_module.h +++ b/include/atframe/modules/worker_pool_module.h @@ -25,7 +25,7 @@ class worker_pool_module : public ::atframework::atapp::module_impl { public: LIBATAPP_MACRO_API worker_pool_module(); - LIBATAPP_MACRO_API virtual ~worker_pool_module(); + LIBATAPP_MACRO_API ~worker_pool_module() override; private: class worker; @@ -78,6 +78,34 @@ class worker_pool_module : public ::atframework::atapp::module_impl { // thread-safe LIBATAPP_MACRO_API std::chrono::microseconds get_tick_interval(const worker_context& context) const; + // thread-safe + LIBATAPP_MACRO_API worker_event_callback_handle_type + add_event_callback_on_worker_created(worker_event_callback_type action); + + // thread-safe + LIBATAPP_MACRO_API void remove_event_callback_on_worker_created(const worker_event_callback_handle_type& handle); + + // thread-safe + LIBATAPP_MACRO_API worker_event_callback_handle_type + add_event_callback_on_worker_removed(worker_event_callback_type action); + + // thread-safe + LIBATAPP_MACRO_API void remove_event_callback_on_worker_removed(const worker_event_callback_handle_type& handle); + + // thread-safe + LIBATAPP_MACRO_API worker_event_callback_handle_type + add_event_callback_on_worker_started(worker_event_callback_type action); + + // thread-safe + LIBATAPP_MACRO_API void remove_event_callback_on_worker_started(const worker_event_callback_handle_type& handle); + + // thread-safe + LIBATAPP_MACRO_API worker_event_callback_handle_type + add_event_callback_on_worker_exiting(worker_event_callback_type action); + + // thread-safe + LIBATAPP_MACRO_API void remove_event_callback_on_worker_exiting(const worker_event_callback_handle_type& handle); + // thread-safe LIBATAPP_MACRO_API size_t get_current_worker_count() const noexcept; diff --git a/src/atframe/modules/worker_pool_module.cpp b/src/atframe/modules/worker_pool_module.cpp index 36c1bff..e43c12d 100644 --- a/src/atframe/modules/worker_pool_module.cpp +++ b/src/atframe/modules/worker_pool_module.cpp @@ -16,7 +16,9 @@ #include #include #include +#include #include +#include #include #include @@ -30,7 +32,14 @@ LIBATAPP_MACRO_NAMESPACE_BEGIN -struct UTIL_SYMBOL_VISIBLE worker_tick_action_handle_data { +struct ATFW_UTIL_SYMBOL_LOCAL worker_event_callback_internal_data { + worker_event_callback_type callback; + std::weak_ptr handle; +}; + +using worker_event_callback_internal_data_pointer = std::shared_ptr; + +struct ATFW_UTIL_SYMBOL_VISIBLE worker_tick_action_handle_data { worker_tick_handle_type type = worker_tick_handle_type::kWorkerTickHandleAny; worker_context specify_worker; size_t version = 0; @@ -38,9 +47,14 @@ struct UTIL_SYMBOL_VISIBLE worker_tick_action_handle_data { const std::list* owner = nullptr; }; +struct ATFW_UTIL_SYMBOL_VISIBLE worker_event_callback_handle_data { + std::list::iterator iter; + std::list* owner = nullptr; +}; + namespace { -struct UTIL_SYMBOL_LOCAL worker_tick_action_container_type { +struct ATFW_UTIL_SYMBOL_LOCAL worker_tick_action_container_type { size_t version = 0; std::list data; }; @@ -53,7 +67,7 @@ enum class worker_status : uint8_t { kExited = 5, }; -struct UTIL_SYMBOL_LOCAL worker_compare_key { +struct ATFW_UTIL_SYMBOL_LOCAL worker_compare_key { size_t pending_job_size; std::chrono::microseconds::rep cpu_time_last_second_busy_us; std::chrono::microseconds::rep cpu_time_last_minute_busy_us; @@ -104,16 +118,64 @@ struct UTIL_SYMBOL_LOCAL worker_compare_key { } }; +static std::atomic& get_worker_unique_id_generator() noexcept { + static std::atomic generator{1}; + return generator; +} + +static worker_event_callback_handle_type internal_add_event_callback( + std::recursive_mutex& lock, std::list& callback_list, + worker_event_callback_type&& action) { + worker_event_callback_handle_type ret = std::make_shared(); + + std::shared_ptr handle_data = + std::make_shared(); + handle_data->callback = std::move(action); + handle_data->handle = ret; + + std::lock_guard lg{lock}; + ret->iter = callback_list.insert(callback_list.end(), std::move(handle_data)); + ret->owner = &callback_list; + + return ret; +} + +static void internal_remove_event_callback(std::recursive_mutex& lock, + std::list& callback_list, + const worker_event_callback_handle_type& handle) { + // 加锁后再检查状态,handle->owner 可能被其他线程的 remove/cleanup 并发修改 + std::lock_guard lg{lock}; + if (handle->owner != &callback_list) { + return; + } + + if (handle->owner->end() == handle->iter) { + handle->owner = nullptr; + return; + } + + handle->owner->erase(handle->iter); + handle->iter = handle->owner->end(); + handle->owner = nullptr; +} + +static std::list collect_event_callback( + std::recursive_mutex& lock, std::list& callback_list) { + std::lock_guard lg{lock}; + return callback_list; +} + } // namespace -class UTIL_SYMBOL_LOCAL worker_pool_module::worker : public std::enable_shared_from_this { +class ATFW_UTIL_SYMBOL_LOCAL worker_pool_module::worker + : public std::enable_shared_from_this { UTIL_DESIGN_PATTERN_NOCOPYABLE(worker); UTIL_DESIGN_PATTERN_NOMOVABLE(worker); friend class worker_pool_module; public: - worker(worker_pool_module::worker_set& owner, uint32_t worker_id); + worker(worker_pool_module::worker_set& owner, uint32_t worker_id, uint64_t worker_unique_id); ~worker(); worker_pool_module::worker_set& get_owner() noexcept { return *owner_; } @@ -214,6 +276,7 @@ class UTIL_SYMBOL_LOCAL worker_pool_module::worker : public std::enable_shared_f ret->specify_worker = get_context(); } else { ret->specify_worker.worker_id = 0; + ret->specify_worker.worker_unique_id = 0; } std::lock_guard lg{tick_handle_lock_}; @@ -293,7 +356,7 @@ class UTIL_SYMBOL_LOCAL worker_pool_module::worker : public std::enable_shared_f std::atomic cpu_time_collect_scaling_down_us_; }; -struct UTIL_SYMBOL_LOCAL worker_pool_module::worker_set { +struct ATFW_UTIL_SYMBOL_LOCAL worker_pool_module::worker_set { bool need_scaling_up; std::atomic closing; std::atomic cleaning; @@ -310,13 +373,22 @@ struct UTIL_SYMBOL_LOCAL worker_pool_module::worker_set { ::tbb::concurrent_queue shared_jobs; + std::recursive_mutex event_lock_on_worker_created; + std::list event_callback_on_worker_created; + std::recursive_mutex event_lock_on_worker_removed; + std::list event_callback_on_worker_removed; + std::recursive_mutex event_lock_on_worker_started; + std::list event_callback_on_worker_started; + std::recursive_mutex event_lock_on_worker_exiting; + std::list event_callback_on_worker_exiting; + worker_set(); UTIL_DESIGN_PATTERN_NOCOPYABLE(worker_set); UTIL_DESIGN_PATTERN_NOMOVABLE(worker_set); }; -struct UTIL_SYMBOL_LOCAL worker_pool_module::scaling_configure { +struct ATFW_UTIL_SYMBOL_LOCAL worker_pool_module::scaling_configure { uint32_t queue_size_limit = 20480; uint32_t max_workers = 4; uint32_t min_workers = 2; @@ -332,7 +404,7 @@ struct UTIL_SYMBOL_LOCAL worker_pool_module::scaling_configure { inline scaling_configure() noexcept {} }; -struct UTIL_SYMBOL_LOCAL worker_pool_module::scaling_statistics { +struct ATFW_UTIL_SYMBOL_LOCAL worker_pool_module::scaling_statistics { std::chrono::system_clock::time_point last_scaling_up_checkpoint; std::chrono::system_clock::time_point last_scaling_down_checkpoint; @@ -344,8 +416,10 @@ struct UTIL_SYMBOL_LOCAL worker_pool_module::scaling_statistics { leak_scan_checkpoint(last_scaling_up_checkpoint) {} }; -worker_pool_module::worker::worker(worker_pool_module::worker_set& owner, uint32_t worker_id) : owner_(&owner) { +worker_pool_module::worker::worker(worker_pool_module::worker_set& owner, uint32_t worker_id, uint64_t worker_unique_id) + : owner_(&owner) { context_.worker_id = worker_id; + context_.worker_unique_id = worker_unique_id; status_.store(static_cast(worker_status::kCreated), std::memory_order_release); created_time_.store(std::chrono::system_clock::now().time_since_epoch().count(), std::memory_order_release); @@ -402,6 +476,16 @@ void worker_pool_module::worker::start(const std::shared_ptr& self, cons self->background_job_thread_ = std::make_shared([self, owner]() { self->status_.store(static_cast(worker_status::kRunning), std::memory_order_release); + { + auto event_on_worker_start = collect_event_callback(self->get_owner().event_lock_on_worker_started, + self->get_owner().event_callback_on_worker_started); + for (auto& fn : event_on_worker_start) { + if (fn && fn->callback) { + fn->callback(self->get_context()); + } + } + } + // loop util end while (!owner->cleaning.load(std::memory_order_acquire)) { if (self->get_context().worker_id > owner->current_expect_workers.load(std::memory_order_acquire) || @@ -516,6 +600,16 @@ void worker_pool_module::worker::start(const std::shared_ptr& self, cons } } + { + auto event_on_worker_exiting = collect_event_callback(self->get_owner().event_lock_on_worker_exiting, + self->get_owner().event_callback_on_worker_exiting); + for (auto& fn : event_on_worker_exiting) { + if (fn && fn->callback) { + fn->callback(self->get_context()); + } + } + } + // exit { std::lock_guard child_lg{self->background_job_lock_}; @@ -937,6 +1031,7 @@ LIBATAPP_MACRO_API int worker_pool_module::spawn(const worker_job_action_pointer *selected_context = worker_ptr->get_context(); } else { selected_context->worker_id = static_cast(worker_type::kMain); + selected_context->worker_unique_id = 0; } } @@ -1054,6 +1149,83 @@ LIBATAPP_MACRO_API std::chrono::microseconds worker_pool_module::get_tick_interv return worker_ptr->get_current_tick_interval(); } +LIBATAPP_MACRO_API worker_event_callback_handle_type +worker_pool_module::add_event_callback_on_worker_created(worker_event_callback_type action) { + if (!worker_set_ || !action) { + return nullptr; + } + return internal_add_event_callback(worker_set_->event_lock_on_worker_created, + worker_set_->event_callback_on_worker_created, std::move(action)); +} + +LIBATAPP_MACRO_API void worker_pool_module::remove_event_callback_on_worker_created( + const worker_event_callback_handle_type& handle) { + if (!handle || !worker_set_) { + return; + } + + internal_remove_event_callback(worker_set_->event_lock_on_worker_created, + worker_set_->event_callback_on_worker_created, handle); +} + +LIBATAPP_MACRO_API worker_event_callback_handle_type +worker_pool_module::add_event_callback_on_worker_removed(worker_event_callback_type action) { + if (!worker_set_ || !action) { + return nullptr; + } + return internal_add_event_callback(worker_set_->event_lock_on_worker_removed, + worker_set_->event_callback_on_worker_removed, std::move(action)); +} + +LIBATAPP_MACRO_API void worker_pool_module::remove_event_callback_on_worker_removed( + const worker_event_callback_handle_type& handle) { + if (!handle || !worker_set_) { + return; + } + + internal_remove_event_callback(worker_set_->event_lock_on_worker_removed, + worker_set_->event_callback_on_worker_removed, handle); +} + +LIBATAPP_MACRO_API worker_event_callback_handle_type +worker_pool_module::add_event_callback_on_worker_started(worker_event_callback_type action) { + if (!worker_set_ || !action) { + return nullptr; + } + + return internal_add_event_callback(worker_set_->event_lock_on_worker_started, + worker_set_->event_callback_on_worker_started, std::move(action)); +} + +LIBATAPP_MACRO_API void worker_pool_module::remove_event_callback_on_worker_started( + const worker_event_callback_handle_type& handle) { + if (!handle || !worker_set_) { + return; + } + + internal_remove_event_callback(worker_set_->event_lock_on_worker_started, + worker_set_->event_callback_on_worker_started, handle); +} + +LIBATAPP_MACRO_API worker_event_callback_handle_type +worker_pool_module::add_event_callback_on_worker_exiting(worker_event_callback_type action) { + if (!worker_set_ || !action) { + return nullptr; + } + return internal_add_event_callback(worker_set_->event_lock_on_worker_exiting, + worker_set_->event_callback_on_worker_exiting, std::move(action)); +} + +LIBATAPP_MACRO_API void worker_pool_module::remove_event_callback_on_worker_exiting( + const worker_event_callback_handle_type& handle) { + if (!handle || !worker_set_) { + return; + } + + internal_remove_event_callback(worker_set_->event_lock_on_worker_exiting, + worker_set_->event_callback_on_worker_exiting, handle); +} + LIBATAPP_MACRO_API size_t worker_pool_module::get_current_worker_count() const noexcept { if (worker_set_) { std::lock_guard lg{worker_set_->worker_lock}; @@ -1072,21 +1244,8 @@ LIBATAPP_MACRO_API void worker_pool_module::foreach_worker_quickly( size_t except_count = get_configure_worker_except_count(); size_t min_count = get_configure_worker_min_count(); - // stable workers always available + uint32_t last_worker_id = 0; bool should_continue = true; - for (uint32_t worker_id = 1; worker_id <= min_count; ++worker_id) { - worker_meta meta = {}; - meta.scaling_mode = worker_scaling_mode::kStable; - should_continue = fn(worker_context{worker_id}, meta); - if (!should_continue) { - break; - } - } - - if (!should_continue) { - return; - } - std::lock_guard lg{worker_set_->worker_lock}; for (auto& worker_ptr : worker_set_->workers) { if (!worker_ptr) { @@ -1099,15 +1258,30 @@ LIBATAPP_MACRO_API void worker_pool_module::foreach_worker_quickly( worker_meta meta = {}; if (worker_ptr->get_context().worker_id <= min_count) { - continue; - } - if (worker_ptr->get_context().worker_id <= except_count) { + meta.scaling_mode = worker_scaling_mode::kStable; + } else if (worker_ptr->get_context().worker_id <= except_count) { meta.scaling_mode = worker_scaling_mode::kDynamic; } else { meta.scaling_mode = worker_scaling_mode::kPendingToDestroy; } + last_worker_id = worker_ptr->get_context().worker_id; if (!fn(worker_ptr->get_context(), meta)) { + should_continue = false; + break; + } + } + + if (!should_continue) { + return; + } + + // stable workers always available + for (uint32_t worker_id = last_worker_id + 1; worker_id <= min_count; ++worker_id) { + worker_meta meta = {}; + meta.scaling_mode = worker_scaling_mode::kStable; + should_continue = fn(worker_context{worker_id, 0}, meta); + if (!should_continue) { break; } } @@ -1133,7 +1307,11 @@ LIBATAPP_MACRO_API void worker_pool_module::foreach_worker( for (uint32_t worker_id = 1; worker_id <= min_count; ++worker_id) { worker_meta meta = {}; meta.scaling_mode = worker_scaling_mode::kStable; - should_continue = fn(worker_context{worker_id}, meta); + uint64_t worker_unique_id = 0; + if (worker_id <= workers.size() && workers[worker_id - 1]) { + worker_unique_id = workers[worker_id - 1]->get_context().worker_unique_id; + } + should_continue = fn(worker_context{worker_id, worker_unique_id}, meta); if (!should_continue) { break; } @@ -1264,6 +1442,7 @@ void worker_pool_module::do_shared_job_on_main_thread() { worker_job_data job_data; worker_context context; context.worker_id = 0; + context.worker_unique_id = 0; std::chrono::system_clock::time_point start_time = std::chrono::system_clock::now(); int32_t no_action_counter = 256; @@ -1297,10 +1476,36 @@ void worker_pool_module::do_scaling_up() { worker_set_->need_scaling_up = false; uint32_t expect_workers = worker_set_->current_expect_workers.load(std::memory_order_acquire); - std::lock_guard lg{worker_set_->worker_lock}; - for (size_t i = worker_set_->workers.size(); i < expect_workers; ++i) { - worker_set_->workers.emplace_back(std::make_shared(*worker_set_, static_cast(i + 1))); - worker::start(worker_set_->workers.back(), worker_set_); + + std::list> new_workers; + { + std::lock_guard lg{worker_set_->worker_lock}; + for (size_t i = worker_set_->workers.size(); i < expect_workers; ++i) { + auto worker_ptr = + std::make_shared(*worker_set_, static_cast(i + 1), + get_worker_unique_id_generator().fetch_add(1, std::memory_order_acq_rel)); + worker_set_->workers.emplace_back(worker_ptr); + new_workers.emplace_back(std::move(worker_ptr)); + } + } + + if (!new_workers.empty()) { + auto event_on_worker_created_data = collect_event_callback(worker_set_->event_lock_on_worker_created, + worker_set_->event_callback_on_worker_created); + + for (auto& worker_ptr : new_workers) { + if (!worker_ptr) { + continue; + } + + for (auto& fn : event_on_worker_created_data) { + if (fn && fn->callback) { + fn->callback(worker_ptr->get_context()); + } + } + + worker::start(worker_ptr, worker_set_); + } } } @@ -1314,20 +1519,42 @@ bool worker_pool_module::internal_reduce_workers() { expect_workers = worker_set_->current_expect_workers.load(std::memory_order_acquire); } - std::lock_guard lg{worker_set_->worker_lock}; - while (worker_set_->workers.size() > expect_workers) { - auto& last_worker = *worker_set_->workers.rbegin(); - if (last_worker) { - if (!last_worker->is_exited()) { - last_worker->wakeup(); - break; + std::list> remove_workers; + + { + std::lock_guard lg{worker_set_->worker_lock}; + while (worker_set_->workers.size() > expect_workers) { + std::shared_ptr last_worker = worker_set_->workers.back(); + if (last_worker) { + if (!last_worker->is_exited()) { + last_worker->wakeup(); + break; + } + + worker_set_->cpu_time_collect_scaling_up_us_for_removed_workers += last_worker->collect_scaling_up_cpu_time(); + worker_set_->cpu_time_collect_scaling_down_us_for_removed_workers += + last_worker->collect_scaling_down_cpu_time(); + remove_workers.emplace_back(std::move(last_worker)); } - worker_set_->cpu_time_collect_scaling_up_us_for_removed_workers += last_worker->collect_scaling_up_cpu_time(); - worker_set_->cpu_time_collect_scaling_down_us_for_removed_workers += last_worker->collect_scaling_down_cpu_time(); + worker_set_->workers.pop_back(); } + } + + if (!remove_workers.empty()) { + auto event_on_worker_removed_data = collect_event_callback(worker_set_->event_lock_on_worker_removed, + worker_set_->event_callback_on_worker_removed); + for (auto& worker_ptr : remove_workers) { + if (!worker_ptr) { + continue; + } - worker_set_->workers.pop_back(); + for (auto& fn : event_on_worker_removed_data) { + if (fn && fn->callback) { + fn->callback(worker_ptr->get_context()); + } + } + } } return !worker_set_->workers.empty(); @@ -1343,45 +1570,70 @@ void worker_pool_module::internal_autofix_workers() { expect_workers = worker_set_->current_expect_workers.load(std::memory_order_acquire); } - std::lock_guard lg{worker_set_->worker_lock}; - bool need_autofix = false; - for (size_t i = 0; !need_autofix && i < worker_set_->workers.size() && i < expect_workers; ++i) { - if (!worker_set_->workers[i]) { - continue; + std::vector> remove_workers; + { + std::lock_guard lg{worker_set_->worker_lock}; + bool need_autofix = false; + for (size_t i = 0; !need_autofix && i < worker_set_->workers.size() && i < expect_workers; ++i) { + if (!worker_set_->workers[i]) { + continue; + } + + if (!worker_set_->workers[i]->is_exited()) { + continue; + } + + need_autofix = true; } - if (!worker_set_->workers[i]->is_exited()) { - continue; + if (!need_autofix) { + return; } - need_autofix = true; - } + std::vector> new_workers; + new_workers.reserve(worker_set_->workers.size()); + for (auto& worker_ptr : worker_set_->workers) { + if (!worker_ptr) { + continue; + } - if (!need_autofix) { - return; - } + if (worker_ptr->is_exited()) { + worker_set_->cpu_time_collect_scaling_up_us_for_removed_workers += worker_ptr->collect_scaling_up_cpu_time(); + worker_set_->cpu_time_collect_scaling_down_us_for_removed_workers += + worker_ptr->collect_scaling_down_cpu_time(); + remove_workers.push_back(worker_ptr); + continue; + } - std::vector> new_workers; - new_workers.reserve(worker_set_->workers.size()); - for (auto& worker_ptr : worker_set_->workers) { - if (!worker_ptr) { - continue; + new_workers.push_back(worker_ptr); } - if (worker_ptr->is_exited()) { - worker_set_->cpu_time_collect_scaling_up_us_for_removed_workers += worker_ptr->collect_scaling_up_cpu_time(); - worker_set_->cpu_time_collect_scaling_down_us_for_removed_workers += worker_ptr->collect_scaling_down_cpu_time(); - continue; + for (size_t i = 0; i < new_workers.size(); ++i) { + new_workers[i]->get_context().worker_id = static_cast(i + 1); } - new_workers.push_back(worker_ptr); + worker_set_->workers.swap(new_workers); } - for (size_t i = 0; i < new_workers.size(); ++i) { - new_workers[i]->get_context().worker_id = static_cast(i + 1); - } + // Trigger event callback for removed workers + if (!remove_workers.empty()) { + auto event_on_worker_removed = collect_event_callback(worker_set_->event_lock_on_worker_removed, + worker_set_->event_callback_on_worker_removed); - worker_set_->workers.swap(new_workers); + if (!event_on_worker_removed.empty()) { + for (auto& worker_ptr : remove_workers) { + if (!worker_ptr) { + continue; + } + + for (auto& fn : event_on_worker_removed) { + if (fn && fn->callback) { + fn->callback(worker_ptr->get_context()); + } + } + } + } + } } void worker_pool_module::internal_cleanup() { @@ -1411,6 +1663,39 @@ void worker_pool_module::internal_cleanup() { worker_set_->configure_tick_max_interval_microseconds.load(std::memory_order_acquire)}; } } + + // cleanup callbacks + std::pair*> callback_list_pair[] = { + {&worker_set_->event_lock_on_worker_created, &worker_set_->event_callback_on_worker_created}, + {&worker_set_->event_lock_on_worker_removed, &worker_set_->event_callback_on_worker_removed}, + {&worker_set_->event_lock_on_worker_started, &worker_set_->event_callback_on_worker_started}, + {&worker_set_->event_lock_on_worker_exiting, &worker_set_->event_callback_on_worker_exiting}}; + + for (auto& callback_lock_and_list : callback_list_pair) { + std::lock_guard lg{*callback_lock_and_list.first}; + + // 先重置所有的handle,以防外部调用remove时访问无效数据 + for (auto& callback_ptr : *callback_lock_and_list.second) { + if (!callback_ptr) { + continue; + } + + if (callback_ptr->handle.expired()) { + continue; + } + + auto handle_ptr = callback_ptr->handle.lock(); + if (!handle_ptr) { + continue; + } + + handle_ptr->owner = nullptr; + handle_ptr->iter = callback_lock_and_list.second->end(); + } + + // 再实际清空数据 + callback_lock_and_list.second->clear(); + } } void worker_pool_module::apply_configure() { diff --git a/test/case/atapp_worker_pool_test.cpp b/test/case/atapp_worker_pool_test.cpp index 7095770..882af23 100644 --- a/test/case/atapp_worker_pool_test.cpp +++ b/test/case/atapp_worker_pool_test.cpp @@ -9,8 +9,10 @@ #include #include #include +#include #include #include +#include #include "frame/test_macros.h" @@ -265,6 +267,206 @@ CASE_TEST(atapp_worker_pool, foreach_stable_workers) { CASE_EXPECT_EQ(min_count + min_count, foreach_counter); } +// event callback +CASE_TEST(atapp_worker_pool, event_callback) { + std::string conf_path_base; + atfw::util::file_system::dirname(__FILE__, 0, conf_path_base); + std::string conf_path = conf_path_base + "/atapp_test_1.yaml"; + + if (!atfw::util::file_system::is_exist(conf_path.c_str())) { + CASE_MSG_INFO() << CASE_MSG_FCOLOR(YELLOW) << conf_path << " not found, skip this test" << std::endl; + return; + } + + atframework::atapp::app app; + const char* args[] = {"app", "-c", conf_path.c_str(), "start"}; + CASE_EXPECT_EQ(0, app.init(nullptr, 4, args, nullptr)); + + auto worker_pool_module = app.get_worker_pool_module(); + CASE_EXPECT_TRUE(!!worker_pool_module); + + if (!worker_pool_module) { + return; + } + + // Invalid input should return empty handle + CASE_EXPECT_FALSE(!!worker_pool_module->add_event_callback_on_worker_created(atapp::worker_event_callback_type{})); + CASE_EXPECT_FALSE(!!worker_pool_module->add_event_callback_on_worker_started(atapp::worker_event_callback_type{})); + CASE_EXPECT_FALSE(!!worker_pool_module->add_event_callback_on_worker_exiting(atapp::worker_event_callback_type{})); + CASE_EXPECT_FALSE(!!worker_pool_module->add_event_callback_on_worker_removed(atapp::worker_event_callback_type{})); + + // Remove an empty handle should be safe + worker_pool_module->remove_event_callback_on_worker_created(atapp::worker_event_callback_handle_type{}); + worker_pool_module->remove_event_callback_on_worker_started(atapp::worker_event_callback_handle_type{}); + worker_pool_module->remove_event_callback_on_worker_exiting(atapp::worker_event_callback_handle_type{}); + worker_pool_module->remove_event_callback_on_worker_removed(atapp::worker_event_callback_handle_type{}); + + // Register and remove before worker created, the removed callback should never be called + std::atomic removed_created_times{0}; + auto removed_created_handle = worker_pool_module->add_event_callback_on_worker_created( + [&removed_created_times](const atapp::worker_context&) { removed_created_times.fetch_add(1); }); + CASE_EXPECT_TRUE(!!removed_created_handle); + worker_pool_module->remove_event_callback_on_worker_created(removed_created_handle); + // Remove an already removed handle should be safe + worker_pool_module->remove_event_callback_on_worker_created(removed_created_handle); + + // Remove a handle by a mismatched event type should be ignored + std::atomic mismatched_created_times{0}; + auto mismatched_created_handle = worker_pool_module->add_event_callback_on_worker_created( + [&mismatched_created_times](const atapp::worker_context&) { mismatched_created_times.fetch_add(1); }); + CASE_EXPECT_TRUE(!!mismatched_created_handle); + worker_pool_module->remove_event_callback_on_worker_started(mismatched_created_handle); + + std::shared_ptr created_lock = std::make_shared(); + std::shared_ptr> created_contexts = + std::make_shared>(); + std::shared_ptr> started_contexts = + std::make_shared>(); + std::shared_ptr> exiting_contexts = + std::make_shared>(); + std::shared_ptr> removed_contexts = + std::make_shared>(); + + auto created_handle = worker_pool_module->add_event_callback_on_worker_created( + [created_lock, created_contexts](const atapp::worker_context& context) { + std::lock_guard lg{*created_lock}; + created_contexts->push_back(context); + }); + auto started_handle = worker_pool_module->add_event_callback_on_worker_started( + [created_lock, started_contexts](const atapp::worker_context& context) { + std::lock_guard lg{*created_lock}; + started_contexts->push_back(context); + }); + auto exiting_handle = worker_pool_module->add_event_callback_on_worker_exiting( + [created_lock, exiting_contexts](const atapp::worker_context& context) { + std::lock_guard lg{*created_lock}; + exiting_contexts->push_back(context); + }); + auto removed_handle = worker_pool_module->add_event_callback_on_worker_removed( + [created_lock, removed_contexts](const atapp::worker_context& context) { + std::lock_guard lg{*created_lock}; + removed_contexts->push_back(context); + }); + + CASE_EXPECT_TRUE(!!created_handle); + CASE_EXPECT_TRUE(!!started_handle); + CASE_EXPECT_TRUE(!!exiting_handle); + CASE_EXPECT_TRUE(!!removed_handle); + + // Workers are created on demand, spawn a job to trigger scaling up + CASE_EXPECT_EQ(0, worker_pool_module->spawn([](const atapp::worker_context&) {})); + + auto min_count = worker_pool_module->get_configure_worker_min_count(); + int32_t sleep_ms = 10000; + while (sleep_ms > 0) { + worker_pool_module->tick(std::chrono::system_clock::now()); + { + std::lock_guard lg{*created_lock}; + if (created_contexts->size() >= min_count && started_contexts->size() >= created_contexts->size()) { + break; + } + } + std::this_thread::sleep_for(std::chrono::milliseconds(20)); + sleep_ms -= 20; + } + + std::vector created_snapshot; + std::vector started_snapshot; + { + std::lock_guard lg{*created_lock}; + created_snapshot = *created_contexts; + started_snapshot = *started_contexts; + } + + // The removed callback should never be called + CASE_EXPECT_EQ(0, removed_created_times.load()); + // The handle removed by a mismatched event type should still be called + CASE_EXPECT_EQ(created_snapshot.size(), mismatched_created_times.load()); + + CASE_EXPECT_GE(created_snapshot.size(), min_count); + CASE_EXPECT_EQ(created_snapshot.size(), started_snapshot.size()); + for (auto& context : created_snapshot) { + CASE_EXPECT_TRUE(atapp::worker_pool_module::is_valid(context)); + CASE_EXPECT_NE(0, context.worker_unique_id); + } + + // worker_unique_id should be unique + for (size_t i = 0; i < created_snapshot.size(); ++i) { + for (size_t j = i + 1; j < created_snapshot.size(); ++j) { + CASE_EXPECT_NE(created_snapshot[i].worker_unique_id, created_snapshot[j].worker_unique_id); + } + } + + // Every created worker should also fire started event with the same worker_unique_id + for (auto& created_context : created_snapshot) { + bool found_started = false; + for (auto& started_context : started_snapshot) { + if (started_context.worker_unique_id == created_context.worker_unique_id) { + found_started = true; + break; + } + } + CASE_EXPECT_TRUE(found_started); + } + + // Stop all workers, exiting and removed events should be fired + worker_pool_module->stop(); + sleep_ms = 10000; + while (sleep_ms > 0) { + worker_pool_module->tick(std::chrono::system_clock::now()); + { + std::lock_guard lg{*created_lock}; + if (removed_contexts->size() >= created_snapshot.size() && exiting_contexts->size() >= created_snapshot.size()) { + break; + } + } + std::this_thread::sleep_for(std::chrono::milliseconds(20)); + sleep_ms -= 20; + } + + std::vector exiting_snapshot; + std::vector removed_snapshot; + { + std::lock_guard lg{*created_lock}; + exiting_snapshot = *exiting_contexts; + removed_snapshot = *removed_contexts; + } + + CASE_EXPECT_GE(exiting_snapshot.size(), created_snapshot.size()); + CASE_EXPECT_GE(removed_snapshot.size(), created_snapshot.size()); + + // Every created worker should also fire exiting and removed events with the same worker_unique_id + for (auto& created_context : created_snapshot) { + bool found_exiting = false; + for (auto& exiting_context : exiting_snapshot) { + if (exiting_context.worker_unique_id == created_context.worker_unique_id) { + found_exiting = true; + break; + } + } + CASE_EXPECT_TRUE(found_exiting); + + bool found_removed = false; + for (auto& removed_context : removed_snapshot) { + if (removed_context.worker_unique_id == created_context.worker_unique_id) { + found_removed = true; + break; + } + } + CASE_EXPECT_TRUE(found_removed); + } + + // Remove callbacks after cleanup should be safe + worker_pool_module->cleanup(); + worker_pool_module->remove_event_callback_on_worker_created(created_handle); + worker_pool_module->remove_event_callback_on_worker_started(started_handle); + worker_pool_module->remove_event_callback_on_worker_exiting(exiting_handle); + worker_pool_module->remove_event_callback_on_worker_removed(removed_handle); + worker_pool_module->remove_event_callback_on_worker_created(mismatched_created_handle); + + CASE_EXPECT_EQ(0, worker_pool_module->stop()); +} + // TODO: spawn with context and ignore the load balance // TODO: scaling up // TODO: scaling down