Network System 0.1.1
High-performance modular networking library for scalable client-server applications
Loading...
Searching...
No Matches
thread_pool_bridge.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
6
7#if KCENON_WITH_COMMON_SYSTEM
9#endif
10
12
14 std::shared_ptr<thread_pool_interface> pool, BackendType backend_type) {
15 if (!pool) {
18 "ThreadPoolBridge requires non-null thread pool",
19 "ThreadPoolBridge::create");
20 }
21 return ok(std::shared_ptr<ThreadPoolBridge>(
22 new ThreadPoolBridge(std::move(pool), backend_type)));
23}
24
26 std::shared_ptr<thread_pool_interface> pool,
27 BackendType backend_type)
28 : pool_(std::move(pool)), backend_type_(backend_type) {
29}
30
32 if (initialized_.load()) {
33 shutdown();
34 }
35}
36
38 if (initialized_.load()) {
39 return error_void(
41 "ThreadPoolBridge already initialized",
42 "ThreadPoolBridge::initialize");
43 }
44
45 if (!pool_) {
46 return error_void(
48 "Thread pool is null",
49 "ThreadPoolBridge::initialize");
50 }
51
52 if (!pool_->is_running()) {
53 return error_void(
55 "Thread pool is not running",
56 "ThreadPoolBridge::initialize");
57 }
58
59 // Check if bridge is enabled (default: true)
60 auto enabled_it = config.properties.find("enabled");
61 if (enabled_it != config.properties.end() && enabled_it->second == "false") {
62 return error_void(
64 "Bridge is disabled in configuration",
65 "ThreadPoolBridge::initialize");
66 }
67
68 // Initialize metrics
69 std::lock_guard<std::mutex> lock(metrics_mutex_);
71 cached_metrics_.last_activity = std::chrono::steady_clock::now();
72 cached_metrics_.custom_metrics["worker_threads"] = static_cast<double>(pool_->worker_count());
73 cached_metrics_.custom_metrics["pending_tasks"] = static_cast<double>(pool_->pending_tasks());
74 cached_metrics_.custom_metrics["backend_type"] = static_cast<double>(backend_type_);
75
76 initialized_.store(true);
77 return ok();
78}
79
81 if (!initialized_.load()) {
82 return ok(); // Idempotent: already shut down
83 }
84
85 // Update metrics to reflect shutdown state
86 {
87 std::lock_guard<std::mutex> lock(metrics_mutex_);
89 cached_metrics_.last_activity = std::chrono::steady_clock::now();
90 }
91
92 initialized_.store(false);
93 return ok();
94}
95
97 return initialized_.load() && pool_ && pool_->is_running();
98}
99
101 std::lock_guard<std::mutex> lock(metrics_mutex_);
102
103 if (!initialized_.load() || !pool_) {
104 BridgeMetrics metrics;
105 metrics.is_healthy = false;
107 return metrics;
108 }
109
110 // Update metrics with current thread pool state
112 metrics.is_healthy = pool_->is_running();
113 metrics.last_activity = std::chrono::steady_clock::now();
114 metrics.custom_metrics["worker_threads"] = static_cast<double>(pool_->worker_count());
115 metrics.custom_metrics["pending_tasks"] = static_cast<double>(pool_->pending_tasks());
116 metrics.custom_metrics["backend_type"] = static_cast<double>(backend_type_);
117
118 // Update cached metrics for next call
119 cached_metrics_ = metrics;
120
121 return metrics;
122}
123
124std::shared_ptr<thread_pool_interface> ThreadPoolBridge::get_thread_pool() const {
125 return pool_;
126}
127
131
133 const std::string& pool_name) {
134 (void)pool_name;
136 if (!pool) {
139 "Failed to get thread pool from thread_integration_manager",
140 "ThreadPoolBridge::from_thread_system");
141 }
142 return create(pool, BackendType::ThreadSystem);
143}
144
145#if KCENON_WITH_COMMON_SYSTEM
146Result<std::shared_ptr<ThreadPoolBridge>> ThreadPoolBridge::from_common_system(
147 std::shared_ptr<::kcenon::common::interfaces::IExecutor> executor) {
148 if (!executor) {
151 "ThreadPoolBridge::from_common_system requires non-null executor",
152 "ThreadPoolBridge::from_common_system");
153 }
154 auto adapter_result = common_to_network_thread_adapter::create(std::move(executor));
155 if (adapter_result.is_err()) {
157 adapter_result.error().code, adapter_result.error().message,
158 "ThreadPoolBridge::from_common_system");
159 }
160 return create(adapter_result.value(), BackendType::CommonSystem);
161}
162#endif
163
164} // namespace kcenon::network::integration
static Result< std::shared_ptr< ThreadPoolBridge > > from_thread_system(const std::string &pool_name="network_pool")
Create bridge from thread_system.
BackendType get_backend_type() const
Get the backend type.
VoidResult shutdown() override
Shutdown the bridge.
static Result< std::shared_ptr< ThreadPoolBridge > > create(std::shared_ptr< thread_pool_interface > pool, BackendType backend_type=BackendType::Custom)
Construct bridge with custom thread pool.
bool is_initialized() const override
Check if the bridge is initialized.
std::shared_ptr< thread_pool_interface > pool_
ThreadPoolBridge(std::shared_ptr< thread_pool_interface > pool, BackendType backend_type)
std::shared_ptr< thread_pool_interface > get_thread_pool() const
Get the underlying thread pool.
BridgeMetrics get_metrics() const override
Get current metrics.
VoidResult initialize(const BridgeConfig &config) override
Initialize the bridge with configuration.
static thread_integration_manager & instance()
Get the singleton instance.
std::shared_ptr< thread_pool_interface > get_thread_pool()
Get the current thread pool.
tracing_config config
Definition exporters.cpp:29
uint32_t code
Definition hpack.cpp:668
VoidResult error_void(int code, const std::string &message, const std::string &source="network_system", const std::string &details="")
VoidResult ok()
Configuration for bridge initialization.
Metrics and health information for a bridge.
std::chrono::steady_clock::time_point last_activity
Timestamp of last activity or health check.
std::map< std::string, double > custom_metrics
Bridge-specific custom metrics.
bool is_healthy
Overall health status of the bridge.
Bidirectional adapters between network_system's thread_pool_interface and common_system's IExecutor/I...
Thread pool integration bridge for network_system.