Network System 0.1.1
High-performance modular networking library for scalable client-server applications
Loading...
Searching...
No Matches
thread_system_adapter.cpp
Go to the documentation of this file.
1// BSD 3-Clause License
2// Copyright (c) 2021-2025, 🍀☀🌕🌥 🌊
3// See the LICENSE file in the project root for full license information.
4
19
21
22#if KCENON_WITH_THREAD_SYSTEM
23
24// Suppress deprecation warnings from thread_system headers
25#pragma clang diagnostic push
26#pragma clang diagnostic ignored "-Wdeprecated-declarations"
27
28#include <thread> // For std::thread::hardware_concurrency and std::this_thread::sleep_for (fallback)
29
30#include <kcenon/thread/core/thread_worker.h>
31
33
34Result<std::shared_ptr<thread_system_pool_adapter>> thread_system_pool_adapter::create(
35 std::shared_ptr<kcenon::thread::thread_pool> pool) {
36 if (!pool) {
39 "thread_system_pool_adapter: pool is null",
40 "thread_system_pool_adapter::create");
41 }
42 return ok(std::shared_ptr<thread_system_pool_adapter>(
43 new thread_system_pool_adapter(std::move(pool))));
44}
45
46thread_system_pool_adapter::thread_system_pool_adapter(
47 std::shared_ptr<kcenon::thread::thread_pool> pool)
48 : pool_(std::move(pool)) {
49}
50
51thread_system_pool_adapter::~thread_system_pool_adapter() {
52 // No scheduler thread to stop - cleanup handled by thread_pool
53}
54
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();
58
59 // Use submit which returns a future and throws on failure
60 try {
61 pool_->submit([task = std::move(task), promise]() mutable {
62 try {
63 if (task) task();
64 promise->set_value();
65 } catch (...) {
66 promise->set_exception(std::current_exception());
67 }
68 });
69 } catch (const std::exception& e) {
70 promise->set_exception(std::make_exception_ptr(
71 std::runtime_error(
72 std::string("thread_system_pool_adapter: submit failed: ") + e.what()
73 )));
74 }
75
76 return future;
77}
78
79std::future<void> thread_system_pool_adapter::submit_delayed(
80 std::function<void()> task,
81 std::chrono::milliseconds delay
82) {
83#if defined(THREAD_HAS_COMMON_EXECUTOR)
84 // Delegate directly to thread_pool::submit_delayed when IExecutor is available
85 // This eliminates the need for a separate scheduler thread
86 return pool_->submit_delayed(std::move(task), delay);
87#else
88 // Fallback: submit a task that sleeps then executes
89 // Note: This blocks a worker thread during the delay period
90 auto promise = std::make_shared<std::promise<void>>();
91 auto future = promise->get_future();
92
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")));
96 return future;
97 }
98
99 // Use submit which returns a future and throws on failure
100 try {
101 pool_->submit([task = std::move(task), delay, promise]() mutable {
102 try {
103 std::this_thread::sleep_for(delay);
104 if (task) task();
105 promise->set_value();
106 } catch (...) {
107 promise->set_exception(std::current_exception());
108 }
109 });
110 } catch (const std::exception& e) {
111 promise->set_exception(std::make_exception_ptr(
112 std::runtime_error(
113 std::string("thread_system_pool_adapter: delayed submit failed: ") + e.what()
114 )));
115 }
116
117 return future;
118#endif
119}
120
121size_t thread_system_pool_adapter::worker_count() const {
122 return pool_->get_active_worker_count();
123}
124
125bool thread_system_pool_adapter::is_running() const {
126 return pool_->is_running();
127}
128
129size_t thread_system_pool_adapter::pending_tasks() const {
130 return pool_->get_pending_task_count();
131}
132
133std::shared_ptr<thread_system_pool_adapter> thread_system_pool_adapter::create_default(
134 const std::string& pool_name
135) {
136 kcenon::thread::thread_context ctx; // default resolves logger/monitoring if registered
137 auto pool = std::make_shared<kcenon::thread::thread_pool>(pool_name, ctx);
138
139 // Add default workers based on hardware concurrency
140 size_t num_threads = std::thread::hardware_concurrency();
141 if (num_threads == 0) {
142 num_threads = 2; // Fallback
143 }
144
145 for (size_t i = 0; i < num_threads; ++i) {
146 pool->enqueue(std::make_unique<kcenon::thread::thread_worker>());
147 }
148
149 (void)pool->start(); // best-effort start; ignore error to keep adapter usable
150 auto result = thread_system_pool_adapter::create(std::move(pool));
151 return result.is_ok() ? result.value() : nullptr;
152}
153
154std::shared_ptr<thread_system_pool_adapter> thread_system_pool_adapter::from_service_or_default(
155 const std::string& pool_name
156) {
157 try {
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();
162 }
163 } catch (...) {
164 // ignore and fallback
165 }
166 return create_default(pool_name);
167}
168
169bool bind_thread_system_pool_into_manager(const std::string& pool_name) {
170 try {
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);
174 return true;
175 } catch (...) {
176 return false;
177 }
178}
179
180} // namespace kcenon::network::integration
181
182#pragma clang diagnostic pop
183
184#endif // KCENON_WITH_THREAD_SYSTEM
185
Feature flags for network_system.
VoidResult ok()
Adapter that bridges thread_system::thread_pool to thread_pool_interface.