1use chrono::{DateTime, Utc};
4use serde::{Deserialize, Serialize};
5#[cfg(all(debug_assertions, feature = "openapi"))]
6use utoipa::ToSchema;
7
8use crate::{Result, TaskCoreError};
9
10#[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 Pending,
17 Active,
19 Succeeded,
21 Failed,
23 Skipped,
25 Canceled,
27}
28
29#[derive(Debug, Clone, Serialize, Deserialize)]
31#[cfg_attr(all(debug_assertions, feature = "openapi"), derive(ToSchema))]
32pub struct TaskStepInfo {
33 pub key: String,
35 pub title: String,
37 pub status: TaskStepStatus,
39 pub progress_current: i64,
41 pub progress_total: i64,
43 pub detail: Option<String>,
45 #[cfg_attr(all(debug_assertions, feature = "openapi"), schema(value_type = Option<String>))]
47 pub started_at: Option<DateTime<Utc>>,
48 #[cfg_attr(all(debug_assertions, feature = "openapi"), schema(value_type = Option<String>))]
50 pub finished_at: Option<DateTime<Utc>>,
51}
52
53#[derive(Debug, Clone, Copy)]
55pub struct TaskStepSpec {
56 pub key: &'static str,
58 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
76pub 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
109pub 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
131pub 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
155pub 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
172pub 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}