Network System 0.1.1
High-performance modular networking library for scalable client-server applications
Loading...
Searching...
No Matches
thread_pool_adapters.h
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
5#pragma once
6
22
23#include <chrono>
24#include <functional>
25#include <future>
26#include <memory>
27#include <stdexcept>
28#include <string>
29
30#if KCENON_WITH_COMMON_SYSTEM
31#include <kcenon/common/interfaces/executor_interface.h>
32#include <kcenon/common/patterns/result.h>
33#endif
34
36
37#if KCENON_WITH_COMMON_SYSTEM
38
46class function_job final : public ::kcenon::common::interfaces::IJob {
47public:
53 explicit function_job(std::function<void()> func, std::string name = "function_job")
54 : func_(std::move(func)), name_(std::move(name)) {}
55
60 ::kcenon::common::VoidResult execute() override {
61 try {
62 if (func_) {
63 func_();
64 }
65 return ::kcenon::common::ok();
66 } catch (const std::exception& ex) {
67 return ::kcenon::common::VoidResult(::kcenon::common::error_info{
68 ::kcenon::common::error_codes::INTERNAL_ERROR,
69 ex.what(),
70 "network_system::function_job"});
71 } catch (...) {
72 return ::kcenon::common::VoidResult(::kcenon::common::error_info{
73 ::kcenon::common::error_codes::INTERNAL_ERROR,
74 "Unknown exception in function_job",
75 "network_system::function_job"});
76 }
77 }
78
83 std::string get_name() const override { return name_; }
84
85private:
86 std::function<void()> func_;
87 std::string name_;
88};
89
107class network_to_common_thread_adapter
108 : public ::kcenon::common::interfaces::IExecutor {
109public:
115 [[nodiscard]] static ::kcenon::common::Result<std::shared_ptr<network_to_common_thread_adapter>>
116 create(std::shared_ptr<thread_pool_interface> pool) {
117 if (!pool) {
118 return ::kcenon::common::Result<std::shared_ptr<network_to_common_thread_adapter>>::err(
119 ::kcenon::common::error_codes::INVALID_ARGUMENT,
120 "network_to_common_thread_adapter requires non-null pool",
121 "network_to_common_thread_adapter::create");
122 }
123 return ::kcenon::common::Result<std::shared_ptr<network_to_common_thread_adapter>>::ok(
124 std::shared_ptr<network_to_common_thread_adapter>(
125 new network_to_common_thread_adapter(std::move(pool))));
126 }
127
133 ::kcenon::common::Result<std::future<void>> execute(
134 std::unique_ptr<::kcenon::common::interfaces::IJob>&& job) override {
135 if (!pool_ || !pool_->is_running()) {
136 return ::kcenon::common::Result<std::future<void>>::err(
137 ::kcenon::common::error_codes::INVALID_ARGUMENT,
138 "Thread pool not running",
139 "network_to_common_thread_adapter");
140 }
141
142 try {
143 auto shared_job = std::shared_ptr<::kcenon::common::interfaces::IJob>(
144 std::move(job));
145 auto prom = std::make_shared<std::promise<void>>();
146 auto fut = prom->get_future();
147
148 pool_->submit([shared_job, prom]() mutable {
149 auto result = shared_job->execute();
150 if (result.is_err()) {
151 prom->set_exception(std::make_exception_ptr(
152 std::runtime_error(result.error().message)));
153 } else {
154 prom->set_value();
155 }
156 });
157
158 return ::kcenon::common::Result<std::future<void>>::ok(std::move(fut));
159 } catch (const std::exception& e) {
160 return ::kcenon::common::Result<std::future<void>>::err(
161 ::kcenon::common::error_codes::INTERNAL_ERROR,
162 e.what(),
163 "network_to_common_thread_adapter");
164 }
165 }
166
173 ::kcenon::common::Result<std::future<void>> execute_delayed(
174 std::unique_ptr<::kcenon::common::interfaces::IJob>&& job,
175 std::chrono::milliseconds delay) override {
176 if (!pool_ || !pool_->is_running()) {
177 return ::kcenon::common::Result<std::future<void>>::err(
178 ::kcenon::common::error_codes::INVALID_ARGUMENT,
179 "Thread pool not running",
180 "network_to_common_thread_adapter");
181 }
182
183 try {
184 auto shared_job = std::shared_ptr<::kcenon::common::interfaces::IJob>(
185 std::move(job));
186 auto prom = std::make_shared<std::promise<void>>();
187 auto fut = prom->get_future();
188
189 pool_->submit_delayed(
190 [shared_job, prom]() mutable {
191 auto result = shared_job->execute();
192 if (result.is_err()) {
193 prom->set_exception(std::make_exception_ptr(
194 std::runtime_error(result.error().message)));
195 } else {
196 prom->set_value();
197 }
198 },
199 delay);
200
201 return ::kcenon::common::Result<std::future<void>>::ok(std::move(fut));
202 } catch (const std::exception& e) {
203 return ::kcenon::common::Result<std::future<void>>::err(
204 ::kcenon::common::error_codes::INTERNAL_ERROR,
205 e.what(),
206 "network_to_common_thread_adapter");
207 }
208 }
209
214 size_t worker_count() const override {
215 return pool_ ? pool_->worker_count() : 0;
216 }
217
222 bool is_running() const override {
223 return pool_ ? pool_->is_running() : false;
224 }
225
230 size_t pending_tasks() const override {
231 return pool_ ? pool_->pending_tasks() : 0;
232 }
233
241 void shutdown([[maybe_unused]] bool wait_for_completion = true) override {
242 // thread_pool_interface doesn't expose shutdown
243 // Lifecycle management should be done on the underlying pool directly
244 }
245
246private:
247 explicit network_to_common_thread_adapter(
248 std::shared_ptr<thread_pool_interface> pool)
249 : pool_(std::move(pool)) {}
250
251 std::shared_ptr<thread_pool_interface> pool_;
252};
253
273class common_to_network_thread_adapter : public thread_pool_interface {
274public:
280 [[nodiscard]] static ::kcenon::common::Result<std::shared_ptr<common_to_network_thread_adapter>>
281 create(std::shared_ptr<::kcenon::common::interfaces::IExecutor> executor) {
282 if (!executor) {
283 return ::kcenon::common::Result<std::shared_ptr<common_to_network_thread_adapter>>::err(
284 ::kcenon::common::error_codes::INVALID_ARGUMENT,
285 "common_to_network_thread_adapter requires non-null executor",
286 "common_to_network_thread_adapter::create");
287 }
288 return ::kcenon::common::Result<std::shared_ptr<common_to_network_thread_adapter>>::ok(
289 std::shared_ptr<common_to_network_thread_adapter>(
290 new common_to_network_thread_adapter(std::move(executor))));
291 }
292
298 std::future<void> submit(std::function<void()> task) override {
299 if (!executor_ || !executor_->is_running()) {
300 return make_error_future("Executor not running");
301 }
302
303 auto result = executor_->execute(
304 std::make_unique<function_job>(std::move(task)));
305
306 if (result.is_err()) {
307 return make_error_future(result.error().message);
308 }
309
310 return std::move(result.value());
311 }
312
319 std::future<void> submit_delayed(
320 std::function<void()> task,
321 std::chrono::milliseconds delay) override {
322 if (!executor_ || !executor_->is_running()) {
323 return make_error_future("Executor not running");
324 }
325
326 auto result = executor_->execute_delayed(
327 std::make_unique<function_job>(std::move(task)), delay);
328
329 if (result.is_err()) {
330 return make_error_future(result.error().message);
331 }
332
333 return std::move(result.value());
334 }
335
340 size_t worker_count() const override {
341 return executor_ ? executor_->worker_count() : 0;
342 }
343
348 bool is_running() const override {
349 return executor_ ? executor_->is_running() : false;
350 }
351
356 size_t pending_tasks() const override {
357 return executor_ ? executor_->pending_tasks() : 0;
358 }
359
364 void shutdown(bool wait_for_completion = true) {
365 if (executor_) {
366 executor_->shutdown(wait_for_completion);
367 }
368 }
369
370private:
371 explicit common_to_network_thread_adapter(
372 std::shared_ptr<::kcenon::common::interfaces::IExecutor> executor)
373 : executor_(std::move(executor)) {}
374
375 static std::future<void> make_error_future(const std::string& message) {
376 std::promise<void> promise;
377 promise.set_exception(std::make_exception_ptr(std::runtime_error(message)));
378 return promise.get_future();
379 }
380
381 std::shared_ptr<::kcenon::common::interfaces::IExecutor> executor_;
382};
383
384#else // !KCENON_WITH_COMMON_SYSTEM
385
386// Placeholder types when common_system is not available
390
394
398
399#endif // KCENON_WITH_COMMON_SYSTEM
400
401} // namespace kcenon::network::integration
Feature flags for network_system.
VoidResult shutdown()
Shutdown the network system.
Thread system integration interface for network_system.