Skip to main content

flare_core_runtime/
config.rs

1//! 运行时配置模块
2//!
3//! 提供运行时配置,包括生命周期、任务启动、健康检查、指标等配置
4
5use std::time::Duration;
6
7/// 通用轮询型后台任务参数(与具体 Broker 无关)
8///
9/// 用于 MQ 消费者、定时轮询任务等
10#[derive(Debug, Clone)]
11pub struct PollWorkerConfig {
12    /// 同时处理中的任务上限(如 Kafka 消费并发)
13    pub concurrency: usize,
14    /// `fetch` 无数据时的休眠间隔
15    pub idle_backoff: Duration,
16    /// `fetch` 失败后的退避间隔
17    pub error_backoff: Duration,
18}
19
20impl Default for PollWorkerConfig {
21    fn default() -> Self {
22        Self {
23            concurrency: 4,
24            idle_backoff: Duration::from_millis(100),
25            error_backoff: Duration::from_secs(1),
26        }
27    }
28}
29
30impl PollWorkerConfig {
31    /// 创建默认配置
32    pub fn new() -> Self {
33        Self::default()
34    }
35
36    /// 设置并发度
37    pub fn with_concurrency(mut self, concurrency: usize) -> Self {
38        self.concurrency = concurrency;
39        self
40    }
41
42    /// 设置空闲退避时间
43    pub fn with_idle_backoff(mut self, d: Duration) -> Self {
44        self.idle_backoff = d;
45        self
46    }
47
48    /// 设置错误退避时间
49    pub fn with_error_backoff(mut self, d: Duration) -> Self {
50        self.error_backoff = d;
51        self
52    }
53}
54
55/// 任务启动配置
56#[derive(Debug, Clone)]
57pub struct TaskStartupConfig {
58    /// 任务启动并发度(默认为 CPU 核心数)
59    pub concurrency: usize,
60    /// 任务启动超时时间(默认 30 秒)
61    pub timeout: Duration,
62    /// 是否启用任务就绪检查(默认 true)
63    pub enable_ready_check: bool,
64    /// 就绪检查超时时间(默认 30 秒)
65    pub ready_check_timeout: Duration,
66}
67
68impl Default for TaskStartupConfig {
69    fn default() -> Self {
70        Self {
71            concurrency: num_cpus::get(),
72            timeout: Duration::from_secs(30),
73            enable_ready_check: true,
74            ready_check_timeout: Duration::from_secs(30),
75        }
76    }
77}
78
79impl TaskStartupConfig {
80    /// 创建默认配置
81    pub fn new() -> Self {
82        Self::default()
83    }
84
85    /// 设置启动并发度
86    pub fn with_concurrency(mut self, concurrency: usize) -> Self {
87        self.concurrency = concurrency;
88        self
89    }
90
91    /// 设置启动超时时间
92    pub fn with_timeout(mut self, timeout: Duration) -> Self {
93        self.timeout = timeout;
94        self
95    }
96
97    /// 启用/禁用就绪检查
98    pub fn with_ready_check(mut self, enable: bool) -> Self {
99        self.enable_ready_check = enable;
100        self
101    }
102
103    /// 设置就绪检查超时时间
104    pub fn with_ready_check_timeout(mut self, timeout: Duration) -> Self {
105        self.ready_check_timeout = timeout;
106        self
107    }
108}
109
110/// 健康检查配置
111#[derive(Debug, Clone)]
112pub struct HealthCheckConfig {
113    /// 健康检查间隔(默认 10 秒)
114    pub interval: Duration,
115    /// 健康检查超时时间(默认 5 秒)
116    pub timeout: Duration,
117    /// 连续失败阈值(默认 3 次)
118    pub failure_threshold: u32,
119    /// 是否启用健康检查(默认 true)
120    pub enabled: bool,
121}
122
123impl Default for HealthCheckConfig {
124    fn default() -> Self {
125        Self {
126            interval: Duration::from_secs(10),
127            timeout: Duration::from_secs(5),
128            failure_threshold: 3,
129            enabled: true,
130        }
131    }
132}
133
134impl HealthCheckConfig {
135    /// 创建默认配置
136    pub fn new() -> Self {
137        Self::default()
138    }
139
140    /// 设置检查间隔
141    pub fn with_interval(mut self, interval: Duration) -> Self {
142        self.interval = interval;
143        self
144    }
145
146    /// 设置超时时间
147    pub fn with_timeout(mut self, timeout: Duration) -> Self {
148        self.timeout = timeout;
149        self
150    }
151
152    /// 设置失败阈值
153    pub fn with_failure_threshold(mut self, threshold: u32) -> Self {
154        self.failure_threshold = threshold;
155        self
156    }
157
158    /// 启用/禁用健康检查
159    pub fn with_enabled(mut self, enabled: bool) -> Self {
160        self.enabled = enabled;
161        self
162    }
163}
164
165/// 指标配置
166#[derive(Debug, Clone)]
167pub struct MetricsConfig {
168    /// 是否启用指标收集(默认 true)
169    pub enabled: bool,
170    /// 指标暴露端口(默认 9090)
171    pub port: u16,
172    /// 指标暴露路径(默认 "/metrics")
173    pub path: String,
174}
175
176impl Default for MetricsConfig {
177    fn default() -> Self {
178        Self {
179            enabled: true,
180            port: 9090,
181            path: "/metrics".to_string(),
182        }
183    }
184}
185
186impl MetricsConfig {
187    /// 创建默认配置
188    pub fn new() -> Self {
189        Self::default()
190    }
191
192    /// 启用/禁用指标收集
193    pub fn with_enabled(mut self, enabled: bool) -> Self {
194        self.enabled = enabled;
195        self
196    }
197
198    /// 设置暴露端口
199    pub fn with_port(mut self, port: u16) -> Self {
200        self.port = port;
201        self
202    }
203
204    /// 设置暴露路径
205    pub fn with_path(mut self, path: impl Into<String>) -> Self {
206        self.path = path.into();
207        self
208    }
209}
210
211/// 微服务运行时配置
212#[derive(Debug, Clone)]
213pub struct RuntimeConfig {
214    /// 关闭超时时间(默认 5 秒)
215    pub shutdown_timeout: Duration,
216    /// 任务启动配置
217    pub task_startup: TaskStartupConfig,
218    /// 健康检查配置
219    pub health_check: HealthCheckConfig,
220    /// 指标配置
221    pub metrics: MetricsConfig,
222    /// 默认轮询工作线程参数
223    pub default_poll_worker: PollWorkerConfig,
224}
225
226impl Default for RuntimeConfig {
227    fn default() -> Self {
228        Self {
229            shutdown_timeout: Duration::from_secs(5),
230            task_startup: TaskStartupConfig::default(),
231            health_check: HealthCheckConfig::default(),
232            metrics: MetricsConfig::default(),
233            default_poll_worker: PollWorkerConfig::default(),
234        }
235    }
236}
237
238impl RuntimeConfig {
239    /// 创建默认配置
240    pub fn new() -> Self {
241        Self::default()
242    }
243
244    /// 设置关闭超时时间
245    pub fn with_shutdown_timeout(mut self, timeout: Duration) -> Self {
246        self.shutdown_timeout = timeout;
247        self
248    }
249
250    /// 设置任务启动配置
251    pub fn with_task_startup(mut self, config: TaskStartupConfig) -> Self {
252        self.task_startup = config;
253        self
254    }
255
256    /// 设置健康检查配置
257    pub fn with_health_check(mut self, config: HealthCheckConfig) -> Self {
258        self.health_check = config;
259        self
260    }
261
262    /// 设置指标配置
263    pub fn with_metrics(mut self, config: MetricsConfig) -> Self {
264        self.metrics = config;
265        self
266    }
267
268    /// 设置默认轮询工作线程参数
269    pub fn with_default_poll_worker(mut self, config: PollWorkerConfig) -> Self {
270        self.default_poll_worker = config;
271        self
272    }
273
274    /// 验证配置
275    pub fn validate(&self) -> Result<(), String> {
276        if self.shutdown_timeout.is_zero() {
277            return Err("shutdown_timeout must be greater than zero".to_string());
278        }
279
280        if self.task_startup.concurrency == 0 {
281            return Err("task_startup.concurrency must be greater than zero".to_string());
282        }
283
284        if self.health_check.failure_threshold == 0 {
285            return Err("health_check.failure_threshold must be greater than zero".to_string());
286        }
287
288        Ok(())
289    }
290}