22#if KCENON_WITH_THREAD_SYSTEM
25#pragma clang diagnostic push
26#pragma clang diagnostic ignored "-Wdeprecated-declarations"
30#include <kcenon/thread/core/thread_worker.h>
34Result<std::shared_ptr<thread_system_pool_adapter>> thread_system_pool_adapter::create(
35 std::shared_ptr<kcenon::thread::thread_pool> pool) {
39 "thread_system_pool_adapter: pool is null",
40 "thread_system_pool_adapter::create");
42 return ok(std::shared_ptr<thread_system_pool_adapter>(
43 new thread_system_pool_adapter(std::move(pool))));
46thread_system_pool_adapter::thread_system_pool_adapter(
47 std::shared_ptr<kcenon::thread::thread_pool> pool)
48 : pool_(std::move(pool)) {
51thread_system_pool_adapter::~thread_system_pool_adapter() {
55std::future<void> thread_system_pool_adapter::submit(std::function<
void()> task) {
56 auto promise = std::make_shared<std::promise<void>>();
57 auto future = promise->get_future();
61 pool_->submit([task = std::move(task), promise]()
mutable {
66 promise->set_exception(std::current_exception());
69 }
catch (
const std::exception& e) {
70 promise->set_exception(std::make_exception_ptr(
72 std::string(
"thread_system_pool_adapter: submit failed: ") + e.what()
79std::future<void> thread_system_pool_adapter::submit_delayed(
80 std::function<
void()> task,
81 std::chrono::milliseconds delay
83#if defined(THREAD_HAS_COMMON_EXECUTOR)
86 return pool_->submit_delayed(std::move(task), delay);
90 auto promise = std::make_shared<std::promise<void>>();
91 auto future = promise->get_future();
93 if (!pool_->is_running()) {
94 promise->set_exception(std::make_exception_ptr(
95 std::runtime_error(
"thread_system_pool_adapter: pool is not running")));
101 pool_->submit([task = std::move(task), delay, promise]()
mutable {
103 std::this_thread::sleep_for(delay);
105 promise->set_value();
107 promise->set_exception(std::current_exception());
110 }
catch (
const std::exception& e) {
111 promise->set_exception(std::make_exception_ptr(
113 std::string(
"thread_system_pool_adapter: delayed submit failed: ") + e.what()
121size_t thread_system_pool_adapter::worker_count()
const {
122 return pool_->get_active_worker_count();
125bool thread_system_pool_adapter::is_running()
const {
126 return pool_->is_running();
129size_t thread_system_pool_adapter::pending_tasks()
const {
130 return pool_->get_pending_task_count();
133std::shared_ptr<thread_system_pool_adapter> thread_system_pool_adapter::create_default(
134 const std::string& pool_name
136 kcenon::thread::thread_context ctx;
137 auto pool = std::make_shared<kcenon::thread::thread_pool>(pool_name, ctx);
140 size_t num_threads = std::thread::hardware_concurrency();
141 if (num_threads == 0) {
145 for (
size_t i = 0; i < num_threads; ++i) {
146 pool->enqueue(std::make_unique<kcenon::thread::thread_worker>());
150 auto result = thread_system_pool_adapter::create(std::move(pool));
151 return result.is_ok() ? result.value() :
nullptr;
154std::shared_ptr<thread_system_pool_adapter> thread_system_pool_adapter::from_service_or_default(
155 const std::string& pool_name
158 auto& sc = kcenon::thread::service_container::global();
159 if (
auto existing = sc.resolve<kcenon::thread::thread_pool>()) {
160 auto result = thread_system_pool_adapter::create(std::move(existing));
161 if (result.is_ok())
return result.value();
166 return create_default(pool_name);
169bool bind_thread_system_pool_into_manager(
const std::string& pool_name) {
171 auto adapter = thread_system_pool_adapter::from_service_or_default(pool_name);
172 if (!adapter)
return false;
173 thread_integration_manager::instance().set_thread_pool(adapter);
182#pragma clang diagnostic pop
Feature flags for network_system.
constexpr int invalid_argument
Adapter that bridges thread_system::thread_pool to thread_pool_interface.