1use 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#[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
28pub 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 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}