Skip to main content

aster_forge_tasks/
steps.rs

1//! Background task step state helpers.
2
3use chrono::{DateTime, Utc};
4use serde::{Deserialize, Serialize};
5#[cfg(all(debug_assertions, feature = "openapi"))]
6use utoipa::ToSchema;
7
8use crate::{Result, TaskCoreError};
9
10/// Runtime status for a task step.
11#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
12#[cfg_attr(all(debug_assertions, feature = "openapi"), derive(ToSchema))]
13#[serde(rename_all = "snake_case")]
14pub enum TaskStepStatus {
15    /// The step has not started.
16    Pending,
17    /// The step is currently running.
18    Active,
19    /// The step completed successfully.
20    Succeeded,
21    /// The step failed.
22    Failed,
23    /// The step was intentionally skipped.
24    Skipped,
25    /// The step was canceled.
26    Canceled,
27}
28
29/// Serialized task step shown in task APIs.
30#[derive(Debug, Clone, Serialize, Deserialize)]
31#[cfg_attr(all(debug_assertions, feature = "openapi"), derive(ToSchema))]
32pub struct TaskStepInfo {
33    /// Stable step key.
34    pub key: String,
35    /// Human-readable step title.
36    pub title: String,
37    /// Current step status.
38    pub status: TaskStepStatus,
39    /// Current progress amount.
40    pub progress_current: i64,
41    /// Total progress amount.
42    pub progress_total: i64,
43    /// Optional detail text.
44    pub detail: Option<String>,
45    /// Step start time.
46    #[cfg_attr(all(debug_assertions, feature = "openapi"), schema(value_type = Option<String>))]
47    pub started_at: Option<DateTime<Utc>>,
48    /// Step finish time.
49    #[cfg_attr(all(debug_assertions, feature = "openapi"), schema(value_type = Option<String>))]
50    pub finished_at: Option<DateTime<Utc>>,
51}
52
53/// Static step definition used to create initial task steps.
54#[derive(Debug, Clone, Copy)]
55pub struct TaskStepSpec {
56    /// Stable step key.
57    pub key: &'static str,
58    /// Human-readable step title.
59    pub title: &'static str,
60}
61
62fn new_task_step(spec: TaskStepSpec, status: TaskStepStatus, detail: Option<&str>) -> TaskStepInfo {
63    let now = (status == TaskStepStatus::Active).then(Utc::now);
64    TaskStepInfo {
65        key: spec.key.to_string(),
66        title: spec.title.to_string(),
67        status,
68        progress_current: 0,
69        progress_total: 0,
70        detail: detail.map(str::to_string),
71        started_at: now,
72        finished_at: None,
73    }
74}
75
76/// Creates initial task step state from static specs.
77pub fn initial_task_steps_from_specs(specs: &[TaskStepSpec]) -> Vec<TaskStepInfo> {
78    specs
79        .iter()
80        .enumerate()
81        .map(|(index, spec)| {
82            new_task_step(
83                *spec,
84                if index == 0 {
85                    TaskStepStatus::Active
86                } else {
87                    TaskStepStatus::Pending
88                },
89                if index == 0 {
90                    Some("Waiting for worker")
91                } else {
92                    None
93                },
94            )
95        })
96        .collect()
97}
98
99fn find_task_step_mut<'a>(
100    steps: &'a mut [TaskStepInfo],
101    key: &str,
102) -> Result<&'a mut TaskStepInfo> {
103    steps
104        .iter_mut()
105        .find(|step| step.key == key)
106        .ok_or_else(|| TaskCoreError::invalid_value(format!("task step '{key}' not found")))
107}
108
109/// Marks a task step active.
110pub fn set_task_step_active(
111    steps: &mut [TaskStepInfo],
112    key: &str,
113    detail: Option<&str>,
114    progress: Option<(i64, i64)>,
115) -> Result<()> {
116    let now = Utc::now();
117    let step = find_task_step_mut(steps, key)?;
118    step.status = TaskStepStatus::Active;
119    if step.started_at.is_none() {
120        step.started_at = Some(now);
121    }
122    step.finished_at = None;
123    step.detail = detail.map(str::to_string);
124    if let Some((current, total)) = progress {
125        step.progress_current = current;
126        step.progress_total = total;
127    }
128    Ok(())
129}
130
131/// Marks a task step succeeded.
132pub fn set_task_step_succeeded(
133    steps: &mut [TaskStepInfo],
134    key: &str,
135    detail: Option<&str>,
136    progress: Option<(i64, i64)>,
137) -> Result<()> {
138    let now = Utc::now();
139    let step = find_task_step_mut(steps, key)?;
140    step.status = TaskStepStatus::Succeeded;
141    if step.started_at.is_none() {
142        step.started_at = Some(now);
143    }
144    step.finished_at = Some(now);
145    step.detail = detail.map(str::to_string);
146    if let Some((current, total)) = progress {
147        step.progress_current = current;
148        step.progress_total = total;
149    } else if step.progress_total > 0 {
150        step.progress_current = step.progress_total;
151    }
152    Ok(())
153}
154
155/// Marks a task step skipped.
156pub fn set_task_step_skipped(
157    steps: &mut [TaskStepInfo],
158    key: &str,
159    detail: Option<&str>,
160) -> Result<()> {
161    let now = Utc::now();
162    let step = find_task_step_mut(steps, key)?;
163    step.status = TaskStepStatus::Skipped;
164    if step.started_at.is_none() {
165        step.started_at = Some(now);
166    }
167    step.finished_at = Some(now);
168    step.detail = detail.map(str::to_string);
169    Ok(())
170}
171
172/// Marks the active step failed, or the last pending step when no step is active.
173pub fn mark_active_step_failed(steps: &mut [TaskStepInfo], detail: Option<&str>) {
174    let now = Utc::now();
175    if let Some(step) = steps
176        .iter_mut()
177        .find(|step| step.status == TaskStepStatus::Active)
178    {
179        step.status = TaskStepStatus::Failed;
180        if step.started_at.is_none() {
181            step.started_at = Some(now);
182        }
183        step.finished_at = Some(now);
184        step.detail = detail.map(str::to_string);
185        return;
186    }
187    if let Some(step) = steps
188        .iter_mut()
189        .rev()
190        .find(|step| step.status == TaskStepStatus::Pending)
191    {
192        step.status = TaskStepStatus::Failed;
193        step.started_at = Some(now);
194        step.finished_at = Some(now);
195        step.detail = detail.map(str::to_string);
196    }
197}
198
199#[cfg(test)]
200mod tests {
201    use super::{
202        TaskStepInfo, TaskStepSpec, TaskStepStatus, initial_task_steps_from_specs,
203        mark_active_step_failed, set_task_step_active, set_task_step_skipped,
204        set_task_step_succeeded,
205    };
206
207    fn step(key: &str, status: TaskStepStatus) -> TaskStepInfo {
208        TaskStepInfo {
209            key: key.to_string(),
210            title: key.to_string(),
211            status,
212            progress_current: 0,
213            progress_total: 1,
214            detail: None,
215            started_at: None,
216            finished_at: None,
217        }
218    }
219
220    #[test]
221    fn initial_steps_activate_first_spec_and_leave_rest_pending() {
222        let steps = initial_task_steps_from_specs(&[
223            TaskStepSpec {
224                key: "prepare",
225                title: "Prepare",
226            },
227            TaskStepSpec {
228                key: "finish",
229                title: "Finish",
230            },
231        ]);
232
233        assert_eq!(steps.len(), 2);
234        assert_eq!(steps[0].key, "prepare");
235        assert_eq!(steps[0].title, "Prepare");
236        assert_eq!(steps[0].status, TaskStepStatus::Active);
237        assert_eq!(steps[0].detail.as_deref(), Some("Waiting for worker"));
238        assert!(steps[0].started_at.is_some());
239        assert_eq!(steps[1].status, TaskStepStatus::Pending);
240        assert_eq!(steps[1].detail, None);
241        assert!(steps[1].started_at.is_none());
242    }
243
244    #[test]
245    fn step_state_helpers_update_timestamps_progress_and_detail() {
246        let mut steps = vec![step("prepare", TaskStepStatus::Pending)];
247
248        set_task_step_active(&mut steps, "prepare", Some("running"), Some((2, 5))).unwrap();
249        assert_eq!(steps[0].status, TaskStepStatus::Active);
250        assert_eq!(steps[0].detail.as_deref(), Some("running"));
251        assert_eq!(steps[0].progress_current, 2);
252        assert_eq!(steps[0].progress_total, 5);
253        assert!(steps[0].started_at.is_some());
254        assert!(steps[0].finished_at.is_none());
255
256        set_task_step_succeeded(&mut steps, "prepare", Some("done"), None).unwrap();
257        assert_eq!(steps[0].status, TaskStepStatus::Succeeded);
258        assert_eq!(steps[0].detail.as_deref(), Some("done"));
259        assert_eq!(steps[0].progress_current, 5);
260        assert!(steps[0].finished_at.is_some());
261
262        set_task_step_skipped(&mut steps, "prepare", Some("skip")).unwrap();
263        assert_eq!(steps[0].status, TaskStepStatus::Skipped);
264        assert_eq!(steps[0].detail.as_deref(), Some("skip"));
265    }
266
267    #[test]
268    fn mark_active_step_failed_updates_active_step_first() {
269        let mut steps = vec![
270            step("prepare", TaskStepStatus::Succeeded),
271            step("process", TaskStepStatus::Active),
272            step("finish", TaskStepStatus::Pending),
273        ];
274
275        mark_active_step_failed(&mut steps, Some("failed"));
276
277        assert_eq!(steps[1].status, TaskStepStatus::Failed);
278        assert_eq!(steps[1].detail.as_deref(), Some("failed"));
279        assert!(steps[1].started_at.is_some());
280        assert!(steps[1].finished_at.is_some());
281        assert_eq!(steps[2].status, TaskStepStatus::Pending);
282    }
283
284    #[test]
285    fn mark_active_step_failed_falls_back_to_last_pending_step() {
286        let mut steps = vec![
287            step("prepare", TaskStepStatus::Succeeded),
288            step("process", TaskStepStatus::Pending),
289            step("finish", TaskStepStatus::Pending),
290        ];
291
292        mark_active_step_failed(&mut steps, Some("pending failed"));
293
294        assert_eq!(steps[1].status, TaskStepStatus::Pending);
295        assert_eq!(steps[2].status, TaskStepStatus::Failed);
296        assert_eq!(steps[2].detail.as_deref(), Some("pending failed"));
297    }
298}