Skip to main content

sz_rust_workflow/engine/
approval.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2024-2026 SZ-Rust Team
3//
4use std::sync::Arc;
5
6use chrono::Utc;
7use uuid::Uuid;
8
9use crate::definition::{FlowDefinition, NodeConfig, NodeType};
10use crate::error::{WorkflowError, WorkflowErrorCode, WorkflowResult};
11use crate::instance::{ApprovalRecord, FlowInstance, InstanceStatus, Task, TaskAction, TaskStatus};
12use crate::observability::{AuditLogger, WorkflowEvent, WorkflowEventBus};
13use crate::repository::{InstanceRepository, TaskRepository};
14use crate::scheduling::approval_strategy::{select_strategy, NodeCompletion};
15use crate::scheduling::candidate_resolver::CandidateResolver;
16
17use super::task_manager::TaskManager;
18
19/// 任务办理结果。
20#[derive(Debug, Clone)]
21pub struct TaskHandleResult {
22    pub task_id: String,
23    pub action: TaskAction,
24    pub node_completion: NodeCompletion,
25    pub advanced: bool,
26}
27
28/// 审批流引擎,对齐 design 2.2.2.4。
29pub struct ApprovalFlowEngine {
30    task_repo: Arc<dyn TaskRepository>,
31    instance_repo: Arc<dyn InstanceRepository>,
32    candidate_resolver: Arc<dyn CandidateResolver>,
33    task_manager: Arc<TaskManager>,
34    audit: Arc<AuditLogger>,
35    event_bus: Arc<dyn WorkflowEventBus>,
36}
37
38impl ApprovalFlowEngine {
39    pub fn new(
40        task_repo: Arc<dyn TaskRepository>,
41        instance_repo: Arc<dyn InstanceRepository>,
42        candidate_resolver: Arc<dyn CandidateResolver>,
43        task_manager: Arc<TaskManager>,
44        audit: Arc<AuditLogger>,
45        event_bus: Arc<dyn WorkflowEventBus>,
46    ) -> Self {
47        Self {
48            task_repo,
49            instance_repo,
50            candidate_resolver,
51            task_manager,
52            audit,
53            event_bus,
54        }
55    }
56
57    /// 办理任务。
58    pub async fn handle(
59        &self,
60        task_id: &str,
61        action: TaskAction,
62        comment: Option<String>,
63        actor: &str,
64        definition: &FlowDefinition,
65    ) -> WorkflowResult<TaskHandleResult> {
66        let task = self.task_repo.get(task_id).await?.ok_or_else(|| {
67            WorkflowError::with_field(
68                WorkflowErrorCode::InstanceNotFound,
69                "任务不存在",
70                "task_id",
71                task_id,
72            )
73        })?;
74
75        if task.status != TaskStatus::Pending {
76            return Err(WorkflowError::with_field(
77                WorkflowErrorCode::TaskNotHandleable,
78                format!("任务状态 {} 不可办理", task.status),
79                "task_id",
80                task_id,
81            ));
82        }
83
84        let instance = self
85            .instance_repo
86            .get(&task.instance_id)
87            .await?
88            .ok_or_else(|| {
89                WorkflowError::with_field(
90                    WorkflowErrorCode::InstanceNotFound,
91                    "实例不存在",
92                    "instance_id",
93                    &task.instance_id,
94                )
95            })?;
96
97        if !instance.status.is_handleable() {
98            return Err(WorkflowError::with_field(
99                WorkflowErrorCode::InstanceNotHandleable,
100                format!("实例状态 {} 不可办理", instance.status),
101                "instance_id",
102                &instance.instance_id,
103            ));
104        }
105
106        if !task.candidates.contains(&actor.to_string()) {
107            return Err(WorkflowError::with_field(
108                WorkflowErrorCode::UnauthorizedHandle,
109                "越权办理:actor 不属于候选人集合",
110                "actor",
111                actor,
112            ));
113        }
114
115        let mut updated_task = task.clone();
116        updated_task.status = match action {
117            TaskAction::Approve => TaskStatus::Completed,
118            TaskAction::Reject => TaskStatus::Rejected,
119            TaskAction::Transfer => TaskStatus::Transferred,
120            _ => TaskStatus::Completed,
121        };
122        updated_task.assignee = Some(actor.into());
123        updated_task.action = Some(action);
124        updated_task.handled_at = Some(Utc::now());
125        self.task_repo.update(&updated_task).await?;
126
127        let record = ApprovalRecord {
128            record_id: Uuid::new_v4().to_string(),
129            instance_id: task.instance_id.clone(),
130            task_id: task.task_id.clone(),
131            node_id: task.node_id.clone(),
132            actor: actor.into(),
133            action,
134            comment,
135            target_user: None,
136            timestamp: Utc::now(),
137        };
138
139        self.event_bus
140            .publish(WorkflowEvent::TaskHandled {
141                instance_id: task.instance_id.clone(),
142                task_id: task.task_id.clone(),
143                actor: actor.into(),
144                action: format!("{:?}", action),
145                timestamp: Utc::now(),
146            })
147            .await
148            .ok();
149
150        self.audit.log_action(
151            crate::observability::AuditAction::Handle,
152            actor,
153            &task.instance_id,
154            instance.status,
155            instance.status,
156            serde_json::to_value(&record).unwrap_or_default(),
157        );
158
159        let node = definition.find_node(&task.node_id).ok_or_else(|| {
160            WorkflowError::with_field(
161                WorkflowErrorCode::InstanceNotFound,
162                "节点不存在",
163                "node_id",
164                &task.node_id,
165            )
166        })?;
167
168        let (strategy_type, next) = match &node.config {
169            NodeConfig::Approval {
170                approval_strategy,
171                next,
172                ..
173            } => (*approval_strategy, next.clone()),
174            _ => {
175                return Ok(TaskHandleResult {
176                    task_id: task_id.into(),
177                    action,
178                    node_completion: NodeCompletion::Completed,
179                    advanced: false,
180                })
181            }
182        };
183
184        let all_tasks = self
185            .task_manager
186            .list_by_instance(&task.instance_id)
187            .await?;
188        let node_tasks: Vec<Task> = all_tasks
189            .iter()
190            .filter(|t| t.node_id == task.node_id)
191            .cloned()
192            .collect();
193        let strategy = select_strategy(strategy_type);
194        let completion = strategy.check_completion(&node_tasks, action);
195
196        let advanced = if completion == NodeCompletion::Completed {
197            self.advance_node(&instance, &task.node_id, &next, definition)
198                .await?;
199            true
200        } else if completion == NodeCompletion::Rejected {
201            self.advance_node(&instance, &task.node_id, &definition.start_node, definition)
202                .await?;
203            true
204        } else {
205            false
206        };
207
208        Ok(TaskHandleResult {
209            task_id: task_id.into(),
210            action,
211            node_completion: completion,
212            advanced,
213        })
214    }
215
216    async fn advance_node(
217        &self,
218        instance: &FlowInstance,
219        current_node: &str,
220        next_node: &str,
221        definition: &FlowDefinition,
222    ) -> WorkflowResult<()> {
223        self.event_bus
224            .publish(WorkflowEvent::NodeLeft {
225                instance_id: instance.instance_id.clone(),
226                node_id: current_node.into(),
227                timestamp: Utc::now(),
228            })
229            .await
230            .ok();
231
232        if let Some(next) = definition.find_node(next_node) {
233            if next.node_type == NodeType::End {
234                let mut updated = instance.clone();
235                updated.status = InstanceStatus::Completed;
236                updated.current_nodes = vec![next_node.into()];
237                updated.bump_version();
238                let _ = self
239                    .instance_repo
240                    .update_with_version(&updated, instance.version_lock)
241                    .await?;
242                self.event_bus
243                    .publish(WorkflowEvent::InstanceCompleted {
244                        instance_id: instance.instance_id.clone(),
245                        timestamp: Utc::now(),
246                    })
247                    .await
248                    .ok();
249            } else if let NodeConfig::Approval {
250                candidate_strategy, ..
251            } = &next.config
252            {
253                let candidates = self
254                    .candidate_resolver
255                    .resolve(candidate_strategy, &instance.context)
256                    .await?;
257                let tasks = self
258                    .task_manager
259                    .create_tasks(&instance.instance_id, next_node, candidates.clone())
260                    .await?;
261                for t in &tasks {
262                    self.event_bus
263                        .publish(WorkflowEvent::TaskCreated {
264                            instance_id: instance.instance_id.clone(),
265                            task_id: t.task_id.clone(),
266                            node_id: next_node.into(),
267                            candidates: candidates.clone(),
268                            timestamp: Utc::now(),
269                        })
270                        .await
271                        .ok();
272                }
273                let mut updated = instance.clone();
274                updated.current_nodes = vec![next_node.into()];
275                updated.bump_version();
276                let _ = self
277                    .instance_repo
278                    .update_with_version(&updated, instance.version_lock)
279                    .await?;
280            }
281            self.event_bus
282                .publish(WorkflowEvent::NodeEntered {
283                    instance_id: instance.instance_id.clone(),
284                    node_id: next_node.into(),
285                    timestamp: Utc::now(),
286                })
287                .await
288                .ok();
289        }
290        Ok(())
291    }
292}