flare_core_runtime/task/spawn.rs
1//! SpawnTask 实现
2//!
3//! 用于包装已经构建好的 Future,是最基础的任务实现
4
5use super::{Task, TaskResult};
6use std::future::Future;
7use std::pin::Pin;
8
9type ShutdownReceiver = tokio::sync::oneshot::Receiver<()>;
10type TaskFuture = Pin<Box<dyn Future<Output = TaskResult> + Send>>;
11type TaskFutureFactory = Box<dyn FnOnce(ShutdownReceiver) -> TaskFuture + Send + 'static>;
12
13/// Spawn 任务
14///
15/// 用于包装已经构建好的 Future(例如 gRPC server)
16/// 用户可以在 service 层构建好 Future,然后通过 runtime 管理
17///
18/// # 特性
19///
20/// - 支持依赖声明
21/// - 支持 shutdown 信号
22/// - 支持优先级设置
23/// - 支持关键任务标记
24///
25/// # 示例
26///
27/// ## 不需要 shutdown 的任务
28///
29/// ```rust
30/// use flare_core_runtime::task::SpawnTask;
31///
32/// let task = SpawnTask::new("my-task", async {
33/// // 任务逻辑
34/// Ok(())
35/// });
36/// ```
37///
38/// ## 需要 shutdown 的任务
39///
40/// ```rust
41/// use flare_core_runtime::task::SpawnTask;
42///
43/// let task = SpawnTask::with_shutdown("my-grpc", |shutdown_rx| {
44/// async move {
45/// // 使用 shutdown_rx 实现优雅停机
46/// let _ = shutdown_rx.await;
47/// Ok(())
48/// }
49/// });
50/// ```
51///
52/// ## 带依赖的任务
53///
54/// ```rust
55/// use flare_core_runtime::task::SpawnTask;
56///
57/// let task = SpawnTask::new("task-b", async { Ok(()) })
58/// .with_dependencies(vec!["task-a".to_string()])
59/// .with_priority(10)
60/// .with_critical(true);
61/// ```
62pub struct SpawnTask {
63 /// 任务名称
64 name: String,
65 /// 任务依赖
66 dependencies: Vec<String>,
67 /// 任务优先级
68 priority: i32,
69 /// 是否关键任务
70 critical: bool,
71 /// Future 构建函数
72 ///
73 /// 使用闭包延迟构建 Future,以便在 run 时传入 shutdown_rx
74 future_fn: TaskFutureFactory,
75}
76
77impl SpawnTask {
78 /// 创建新的 spawn 任务(不需要 shutdown_rx)
79 ///
80 /// # 参数
81 ///
82 /// * `name` - 任务名称
83 /// * `future` - 要运行的 Future(不依赖 shutdown_rx)
84 ///
85 /// # 示例
86 ///
87 /// ```rust
88 /// use flare_core_runtime::task::SpawnTask;
89 ///
90 /// let task = SpawnTask::new("my-task", async {
91 /// // 任务逻辑
92 /// Ok(())
93 /// });
94 /// ```
95 pub fn new<Fut>(name: impl Into<String>, future: Fut) -> Self
96 where
97 Fut: Future<Output = TaskResult> + Send + 'static,
98 {
99 Self {
100 name: name.into(),
101 dependencies: Vec::new(),
102 priority: 0,
103 critical: false,
104 future_fn: Box::new(move |_shutdown_rx| Box::pin(future)),
105 }
106 }
107
108 /// 创建新的 spawn 任务(需要 shutdown_rx)
109 ///
110 /// # 参数
111 ///
112 /// * `name` - 任务名称
113 /// * `future_fn` - 闭包,接收 shutdown_rx,返回 Future
114 ///
115 /// # 示例
116 ///
117 /// ```rust
118 /// use flare_core_runtime::task::SpawnTask;
119 ///
120 /// let task = SpawnTask::with_shutdown("my-task", |shutdown_rx| {
121 /// async move {
122 /// // 使用 shutdown_rx
123 /// let _ = shutdown_rx.await;
124 /// Ok(())
125 /// }
126 /// });
127 /// ```
128 pub fn with_shutdown<F, Fut>(name: impl Into<String>, future_fn: F) -> Self
129 where
130 F: FnOnce(tokio::sync::oneshot::Receiver<()>) -> Fut + Send + 'static,
131 Fut: Future<Output = TaskResult> + Send + 'static,
132 {
133 Self {
134 name: name.into(),
135 dependencies: Vec::new(),
136 priority: 0,
137 critical: false,
138 future_fn: Box::new(move |shutdown_rx| Box::pin(future_fn(shutdown_rx))),
139 }
140 }
141
142 /// 设置任务依赖
143 ///
144 /// # 参数
145 ///
146 /// * `deps` - 依赖的任务名称列表
147 ///
148 /// # 示例
149 ///
150 /// ```rust
151 /// use flare_core_runtime::task::SpawnTask;
152 ///
153 /// let task = SpawnTask::new("task-b", async { Ok(()) })
154 /// .with_dependencies(vec!["task-a".to_string()]);
155 /// ```
156 pub fn with_dependencies(mut self, deps: Vec<String>) -> Self {
157 self.dependencies = deps;
158 self
159 }
160
161 /// 设置任务优先级
162 ///
163 /// # 参数
164 ///
165 /// * `priority` - 优先级(数值越大优先级越高)
166 ///
167 /// # 示例
168 ///
169 /// ```rust
170 /// use flare_core_runtime::task::SpawnTask;
171 ///
172 /// let task = SpawnTask::new("my-task", async { Ok(()) })
173 /// .with_priority(10);
174 /// ```
175 pub fn with_priority(mut self, priority: i32) -> Self {
176 self.priority = priority;
177 self
178 }
179
180 /// 设置是否为关键任务
181 ///
182 /// # 参数
183 ///
184 /// * `critical` - 是否为关键任务
185 ///
186 /// # 示例
187 ///
188 /// ```rust
189 /// use flare_core_runtime::task::SpawnTask;
190 ///
191 /// let task = SpawnTask::new("my-task", async { Ok(()) })
192 /// .with_critical(true);
193 /// ```
194 pub fn with_critical(mut self, critical: bool) -> Self {
195 self.critical = critical;
196 self
197 }
198}
199
200impl Task for SpawnTask {
201 fn name(&self) -> &str {
202 &self.name
203 }
204
205 fn dependencies(&self) -> Vec<String> {
206 self.dependencies.clone()
207 }
208
209 fn run(
210 self: Box<Self>,
211 shutdown_rx: tokio::sync::oneshot::Receiver<()>,
212 ) -> Pin<Box<dyn Future<Output = TaskResult> + Send>> {
213 // 调用 future_fn,传入 shutdown_rx
214 (self.future_fn)(shutdown_rx)
215 }
216
217 fn priority(&self) -> i32 {
218 self.priority
219 }
220
221 fn is_critical(&self) -> bool {
222 self.critical
223 }
224}
225
226#[cfg(test)]
227mod tests {
228 use super::*;
229 use tokio::sync::oneshot;
230
231 #[tokio::test]
232 async fn test_spawn_task_new() {
233 let task = SpawnTask::new("test-task", async { Ok(()) });
234 assert_eq!(task.name(), "test-task");
235 assert!(task.dependencies().is_empty());
236 assert_eq!(task.priority(), 0);
237 assert!(!task.is_critical());
238 }
239
240 #[tokio::test]
241 async fn test_spawn_task_with_dependencies() {
242 let task = SpawnTask::new("test-task", async { Ok(()) })
243 .with_dependencies(vec!["dep-1".to_string(), "dep-2".to_string()]);
244
245 assert_eq!(task.dependencies(), vec!["dep-1", "dep-2"]);
246 }
247
248 #[tokio::test]
249 async fn test_spawn_task_with_priority() {
250 let task = SpawnTask::new("test-task", async { Ok(()) }).with_priority(10);
251
252 assert_eq!(task.priority(), 10);
253 }
254
255 #[tokio::test]
256 async fn test_spawn_task_with_critical() {
257 let task = SpawnTask::new("test-task", async { Ok(()) }).with_critical(true);
258
259 assert!(task.is_critical());
260 }
261
262 #[tokio::test]
263 async fn test_spawn_task_run() {
264 let task = SpawnTask::new("test-task", async { Ok(()) });
265 let (_tx, rx) = oneshot::channel();
266
267 let result = Box::new(task).run(rx).await;
268 assert!(result.is_ok());
269 }
270
271 #[tokio::test]
272 async fn test_spawn_task_with_shutdown_run() {
273 let task = SpawnTask::with_shutdown("test-task", |shutdown_rx| {
274 async move {
275 // 模拟等待 shutdown 信号
276 tokio::time::timeout(std::time::Duration::from_millis(100), shutdown_rx)
277 .await
278 .ok();
279 Ok(())
280 }
281 });
282
283 let (tx, rx) = oneshot::channel();
284
285 // 发送 shutdown 信号
286 tx.send(()).unwrap();
287
288 let result = Box::new(task).run(rx).await;
289 assert!(result.is_ok());
290 }
291}