Skip to main content

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}