Skip to main content

sz_rust_core/runtime/
worker.rs

1//! Worker 数量配置
2//!
3//! ## PHP 对齐
4//!
5//! 对齐 PHP `think-swoole` 的 `worker_num` 配置项:
6//!
7//! ```php
8//! // vendor/topthink/think-swoole/src/config/swoole.php:17
9//! 'options' => [
10//!     'reactor_num'     => swoole_cpu_num(),
11//!     'worker_num'      => swoole_cpu_num(),   // ← 默认 = CPU 核数
12//!     'task_worker_num' => swoole_cpu_num(),
13//! ]
14//! ```
15//!
16//! Rust 端使用 `num_cpus::get()` 获取 CPU 核数,并提供 builder 模式自定义。
17
18use std::fmt;
19
20/// 默认 worker 数(对齐 `swoole_cpu_num()`)
21pub const DEFAULT_WORKER_NUM: usize = 0; // 0 表示使用 num_cpus::get()
22
23/// 最小 worker 数(避免 0 worker 导致 runtime 无法启动)
24pub const MIN_WORKER_NUM: usize = 1;
25
26/// 最大 worker 数(避免过度创建线程导致调度开销过大)
27pub const MAX_WORKER_NUM: usize = 256;
28
29/// Worker 数量配置
30///
31/// 对齐 PHP `swoole.php` 配置中的 `worker_num` / `reactor_num` / `task_worker_num`。
32///
33/// ## 字段
34///
35/// - `worker_num`:业务 worker 数(处理 HTTP 请求等),默认 = CPU 核数
36/// - `reactor_num`:reactor 线程数(对齐 Swoole reactor,Rust 端保留为元数据)
37/// - `task_worker_num`:task worker 数(对齐 Swoole task worker,用于异步任务)
38///
39/// ## 用法
40///
41/// ```rust,ignore
42/// use sz_rust_core::runtime::worker::WorkerConfig;
43///
44/// let config = WorkerConfig::new();
45/// assert_eq!(config.worker_num(), num_cpus::get());
46///
47/// let custom = WorkerConfig::new().with_worker_num(8);
48/// assert_eq!(custom.worker_num(), 8);
49/// ```
50#[derive(Debug, Clone, PartialEq, Eq)]
51pub struct WorkerConfig {
52    /// 业务 worker 数(0 表示使用 CPU 核数)
53    worker_num: usize,
54    /// reactor 线程数(保留字段,对齐 PHP)
55    reactor_num: usize,
56    /// task worker 数(保留字段,对齐 PHP)
57    task_worker_num: usize,
58}
59
60impl WorkerConfig {
61    /// 创建默认配置:所有字段 = `num_cpus::get()`(对齐 `swoole_cpu_num()`)
62    pub fn new() -> Self {
63        let cpu = num_cpus::get();
64        Self {
65            worker_num: cpu,
66            reactor_num: cpu,
67            task_worker_num: cpu,
68        }
69    }
70
71    /// 自定义 worker_num(对齐 `worker_num` 配置项)
72    ///
73    /// - `n = 0` 会被强制为 `num_cpus::get()`
74    /// - `n > MAX_WORKER_NUM` 会被截断为 `MAX_WORKER_NUM`
75    pub fn with_worker_num(mut self, n: usize) -> Self {
76        self.worker_num = if n == 0 {
77            num_cpus::get()
78        } else {
79            n.clamp(MIN_WORKER_NUM, MAX_WORKER_NUM)
80        };
81        self
82    }
83
84    /// 自定义 reactor_num(对齐 `reactor_num` 配置项)
85    pub fn with_reactor_num(mut self, n: usize) -> Self {
86        self.reactor_num = if n == 0 {
87            num_cpus::get()
88        } else {
89            n.clamp(MIN_WORKER_NUM, MAX_WORKER_NUM)
90        };
91        self
92    }
93
94    /// 自定义 task_worker_num(对齐 `task_worker_num` 配置项)
95    pub fn with_task_worker_num(mut self, n: usize) -> Self {
96        self.task_worker_num = if n == 0 {
97            num_cpus::get()
98        } else {
99            n.clamp(MIN_WORKER_NUM, MAX_WORKER_NUM)
100        };
101        self
102    }
103
104    /// 获取 worker_num
105    pub fn worker_num(&self) -> usize {
106        self.worker_num
107    }
108
109    /// 获取 reactor_num
110    pub fn reactor_num(&self) -> usize {
111        self.reactor_num
112    }
113
114    /// 获取 task_worker_num
115    pub fn task_worker_num(&self) -> usize {
116        self.task_worker_num
117    }
118
119    /// 获取 CPU 核数(对齐 `swoole_cpu_num()`)
120    pub fn cpu_num(&self) -> usize {
121        num_cpus::get()
122    }
123
124    /// 验证配置是否有效
125    pub fn validate(&self) -> bool {
126        self.worker_num >= MIN_WORKER_NUM
127            && self.worker_num <= MAX_WORKER_NUM
128            && self.reactor_num >= MIN_WORKER_NUM
129            && self.reactor_num <= MAX_WORKER_NUM
130            && self.task_worker_num >= MIN_WORKER_NUM
131            && self.task_worker_num <= MAX_WORKER_NUM
132    }
133
134    /// 从环境变量读取 worker_num(对齐 PHP `env('SWOOLE_WORKER_NUM')`)
135    ///
136    /// - 环境变量 `SZ_RUST_WORKER_NUM` 优先级最高
137    /// - 未设置或解析失败时使用默认值(CPU 核数)
138    pub fn from_env() -> Self {
139        let mut config = Self::new();
140        if let Ok(val) = std::env::var("SZ_RUST_WORKER_NUM") {
141            if let Ok(n) = val.parse::<usize>() {
142                config = config.with_worker_num(n);
143            }
144        }
145        config
146    }
147}
148
149impl Default for WorkerConfig {
150    fn default() -> Self {
151        Self::new()
152    }
153}
154
155impl fmt::Display for WorkerConfig {
156    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
157        write!(
158            f,
159            "WorkerConfig{{worker_num={}, reactor_num={}, task_worker_num={}, cpu={}}}",
160            self.worker_num,
161            self.reactor_num,
162            self.task_worker_num,
163            self.cpu_num()
164        )
165    }
166}
167
168#[cfg(test)]
169mod tests {
170    use super::*;
171
172    /// env 测试全局互斥锁:`std::env::set_var/remove_var` 非线程安全,
173    /// 并行测试共享进程环境变量会互相污染(P3 竞态修复:test_from_env_* 串行化)
174    static ENV_TEST_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
175
176    #[test]
177    fn test_new_defaults_to_cpu_count() {
178        let config = WorkerConfig::new();
179        let cpu = num_cpus::get();
180        assert_eq!(config.worker_num(), cpu);
181        assert_eq!(config.reactor_num(), cpu);
182        assert_eq!(config.task_worker_num(), cpu);
183    }
184
185    #[test]
186    fn test_with_worker_num_custom() {
187        let config = WorkerConfig::new().with_worker_num(8);
188        assert_eq!(config.worker_num(), 8);
189    }
190
191    #[test]
192    fn test_with_worker_num_zero_falls_back_to_cpu() {
193        let config = WorkerConfig::new().with_worker_num(0);
194        assert_eq!(config.worker_num(), num_cpus::get());
195    }
196
197    #[test]
198    fn test_with_worker_num_exceeds_max_clamped() {
199        let config = WorkerConfig::new().with_worker_num(1024);
200        assert_eq!(config.worker_num(), MAX_WORKER_NUM);
201    }
202
203    #[test]
204    fn test_with_reactor_num_custom() {
205        let config = WorkerConfig::new().with_reactor_num(4);
206        assert_eq!(config.reactor_num(), 4);
207    }
208
209    #[test]
210    fn test_with_task_worker_num_custom() {
211        let config = WorkerConfig::new().with_task_worker_num(16);
212        assert_eq!(config.task_worker_num(), 16);
213    }
214
215    #[test]
216    fn test_cpu_num_matches_num_cpus() {
217        let config = WorkerConfig::new();
218        assert_eq!(config.cpu_num(), num_cpus::get());
219    }
220
221    #[test]
222    fn test_validate_valid_config() {
223        let config = WorkerConfig::new();
224        assert!(config.validate());
225    }
226
227    #[test]
228    fn test_validate_boundary_values() {
229        let config_min = WorkerConfig::new()
230            .with_worker_num(MIN_WORKER_NUM)
231            .with_reactor_num(MIN_WORKER_NUM)
232            .with_task_worker_num(MIN_WORKER_NUM);
233        assert!(config_min.validate());
234
235        let config_max = WorkerConfig::new()
236            .with_worker_num(MAX_WORKER_NUM)
237            .with_reactor_num(MAX_WORKER_NUM)
238            .with_task_worker_num(MAX_WORKER_NUM);
239        assert!(config_max.validate());
240    }
241
242    #[test]
243    fn test_from_env_default() {
244        // 清除环境变量确保使用默认值
245        // 注意:env 读写非线程安全,多个 env 测试并行时需互斥串行(P3 竞态修复)
246        let _guard = ENV_TEST_LOCK.lock();
247        std::env::remove_var("SZ_RUST_WORKER_NUM");
248        let config = WorkerConfig::from_env();
249        assert_eq!(config.worker_num(), num_cpus::get());
250    }
251
252    #[test]
253    fn test_from_env_custom() {
254        let _guard = ENV_TEST_LOCK.lock();
255        std::env::set_var("SZ_RUST_WORKER_NUM", "12");
256        let config = WorkerConfig::from_env();
257        assert_eq!(config.worker_num(), 12);
258        std::env::remove_var("SZ_RUST_WORKER_NUM");
259    }
260
261    #[test]
262    fn test_from_env_invalid_falls_back_to_default() {
263        let _guard = ENV_TEST_LOCK.lock();
264        std::env::set_var("SZ_RUST_WORKER_NUM", "not-a-number");
265        let config = WorkerConfig::from_env();
266        assert_eq!(config.worker_num(), num_cpus::get());
267        std::env::remove_var("SZ_RUST_WORKER_NUM");
268    }
269
270    #[test]
271    fn test_display_format() {
272        let config = WorkerConfig::new().with_worker_num(4);
273        let s = format!("{}", config);
274        assert!(s.contains("worker_num=4"));
275        assert!(s.contains("reactor_num="));
276        assert!(s.contains("task_worker_num="));
277    }
278
279    #[test]
280    fn test_default_equals_new() {
281        let config1 = WorkerConfig::default();
282        let config2 = WorkerConfig::new();
283        assert_eq!(config1, config2);
284    }
285
286    #[test]
287    fn test_clone_and_equality() {
288        let config1 = WorkerConfig::new().with_worker_num(4);
289        let config2 = config1.clone();
290        assert_eq!(config1, config2);
291    }
292
293    #[test]
294    fn test_builder_chaining() {
295        let config = WorkerConfig::new()
296            .with_worker_num(4)
297            .with_reactor_num(2)
298            .with_task_worker_num(8);
299        assert_eq!(config.worker_num(), 4);
300        assert_eq!(config.reactor_num(), 2);
301        assert_eq!(config.task_worker_num(), 8);
302        assert!(config.validate());
303    }
304}