30#if KCENON_WITH_COMMON_SYSTEM
31#include <kcenon/common/interfaces/executor_interface.h>
32#include <kcenon/common/patterns/result.h>
37#if KCENON_WITH_COMMON_SYSTEM
46class function_job final :
public ::kcenon::common::interfaces::IJob {
53 explicit function_job(std::function<
void()> func, std::string name =
"function_job")
54 : func_(std::move(func)), name_(std::move(name)) {}
60 ::kcenon::common::VoidResult execute()
override {
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,
70 "network_system::function_job"});
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"});
83 std::string get_name()
const override {
return name_; }
86 std::function<void()> func_;
107class network_to_common_thread_adapter
108 :
public ::kcenon::common::interfaces::IExecutor {
115 [[nodiscard]] static ::kcenon::common::Result<std::shared_ptr<network_to_common_thread_adapter>>
116 create(std::shared_ptr<thread_pool_interface> 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");
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))));
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");
143 auto shared_job = std::shared_ptr<::kcenon::common::interfaces::IJob>(
145 auto prom = std::make_shared<std::promise<void>>();
146 auto fut = prom->get_future();
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)));
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,
163 "network_to_common_thread_adapter");
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");
184 auto shared_job = std::shared_ptr<::kcenon::common::interfaces::IJob>(
186 auto prom = std::make_shared<std::promise<void>>();
187 auto fut = prom->get_future();
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)));
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,
206 "network_to_common_thread_adapter");
214 size_t worker_count()
const override {
215 return pool_ ? pool_->worker_count() : 0;
222 bool is_running()
const override {
223 return pool_ ? pool_->is_running() :
false;
230 size_t pending_tasks()
const override {
231 return pool_ ? pool_->pending_tasks() : 0;
241 void shutdown([[maybe_unused]]
bool wait_for_completion =
true)
override {
247 explicit network_to_common_thread_adapter(
248 std::shared_ptr<thread_pool_interface> pool)
249 : pool_(std::move(pool)) {}
251 std::shared_ptr<thread_pool_interface> pool_;
273class common_to_network_thread_adapter :
public thread_pool_interface {
280 [[nodiscard]] static ::kcenon::common::Result<std::shared_ptr<common_to_network_thread_adapter>>
281 create(std::shared_ptr<::kcenon::common::interfaces::IExecutor> 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");
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))));
298 std::future<void> submit(std::function<
void()> task)
override {
299 if (!executor_ || !executor_->is_running()) {
300 return make_error_future(
"Executor not running");
303 auto result = executor_->execute(
304 std::make_unique<function_job>(std::move(task)));
306 if (result.is_err()) {
307 return make_error_future(result.error().message);
310 return std::move(result.value());
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");
326 auto result = executor_->execute_delayed(
327 std::make_unique<function_job>(std::move(task)), delay);
329 if (result.is_err()) {
330 return make_error_future(result.error().message);
333 return std::move(result.value());
340 size_t worker_count()
const override {
341 return executor_ ? executor_->worker_count() : 0;
348 bool is_running()
const override {
349 return executor_ ? executor_->is_running() :
false;
356 size_t pending_tasks()
const override {
357 return executor_ ? executor_->pending_tasks() : 0;
364 void shutdown(
bool wait_for_completion =
true) {
366 executor_->shutdown(wait_for_completion);
371 explicit common_to_network_thread_adapter(
372 std::shared_ptr<::kcenon::common::interfaces::IExecutor> executor)
373 : executor_(std::move(executor)) {}
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();
381 std::shared_ptr<::kcenon::common::interfaces::IExecutor> executor_;
Feature flags for network_system.
VoidResult shutdown()
Shutdown the network system.
common_to_network_thread_adapter_unavailable()=delete
function_job_unavailable()=delete
network_to_common_thread_adapter_unavailable()=delete
Thread system integration interface for network_system.