1use std::sync::Arc;
5
6use crate::api::{DesignerApi, VersionManager};
7use crate::config::WorkflowConfig;
8use crate::definition::{DefinitionFormat, DefinitionValidator, ValidationIssue};
9use crate::deps::WorkflowDeps;
10use crate::engine::approval::{ApprovalFlowEngine, TaskHandleResult};
11use crate::engine::history::HistoryRecorder;
12use crate::engine::instance::{InstanceDetail, InstanceManager, InstanceSummary};
13use crate::engine::state_machine::{StateMachineEngine, TransitionResult};
14use crate::engine::task_manager::TaskManager;
15use crate::error::WorkflowResult;
16use crate::instance::{HistoryEntry, PageRequest, PageResult, Task, TaskAction};
17use crate::repository::DefinitionId;
18
19pub struct WorkflowEngine {
21 #[allow(dead_code)]
22 config: WorkflowConfig,
23 designer_api: Arc<DesignerApi>,
24 instance_manager: Arc<InstanceManager>,
25 state_machine: Arc<StateMachineEngine>,
26 approval_flow: Arc<ApprovalFlowEngine>,
27 task_manager: Arc<TaskManager>,
28 #[allow(dead_code)] history_recorder: Arc<HistoryRecorder>,
30}
31
32impl WorkflowEngine {
33 pub fn new(config: WorkflowConfig, deps: WorkflowDeps) -> Self {
34 let task_manager = Arc::new(TaskManager::new(deps.task_repo.clone()));
35 let history_recorder = Arc::new(HistoryRecorder::new(
36 deps.history_repo.clone(),
37 deps.sensitive_registry.clone(),
38 ));
39 let version_manager = Arc::new(VersionManager::new(deps.definition_repo.clone()));
40 let designer_api = Arc::new(DesignerApi::new(
41 DefinitionValidator::new(Arc::new(crate::definition::NoopPluginChecker)),
42 deps.definition_repo.clone(),
43 version_manager,
44 ));
45 let instance_manager = Arc::new(InstanceManager::new(
46 deps.instance_repo.clone(),
47 deps.definition_repo.clone(),
48 task_manager.clone(),
49 deps.candidate_resolver.clone(),
50 history_recorder.clone(),
51 deps.audit.clone(),
52 deps.event_bus.clone(),
53 deps.sensitive_registry.clone(),
54 ));
55 let state_machine = Arc::new(StateMachineEngine::new(
56 deps.instance_repo.clone(),
57 deps.guard_evaluator.clone(),
58 deps.event_bus.clone(),
59 deps.audit.clone(),
60 ));
61 let approval_flow = Arc::new(ApprovalFlowEngine::new(
62 deps.task_repo.clone(),
63 deps.instance_repo.clone(),
64 deps.candidate_resolver.clone(),
65 task_manager.clone(),
66 deps.audit.clone(),
67 deps.event_bus.clone(),
68 ));
69
70 Self {
71 config,
72 designer_api,
73 instance_manager,
74 state_machine,
75 approval_flow,
76 task_manager,
77 history_recorder,
78 }
79 }
80
81 pub async fn validate_definition(
82 &self,
83 text: &str,
84 format: DefinitionFormat,
85 ) -> WorkflowResult<Vec<ValidationIssue>> {
86 self.designer_api.validate_definition(text, format).await
87 }
88
89 pub async fn import_definition(
90 &self,
91 text: &str,
92 format: DefinitionFormat,
93 ) -> WorkflowResult<DefinitionId> {
94 self.designer_api.import_definition(text, format).await
95 }
96
97 pub async fn export_definition(
98 &self,
99 id: &DefinitionId,
100 format: DefinitionFormat,
101 ) -> WorkflowResult<String> {
102 self.designer_api.export_definition(id, format).await
103 }
104
105 pub async fn start_instance(
106 &self,
107 flow_key: &str,
108 context: serde_json::Value,
109 initiator: &str,
110 ) -> WorkflowResult<InstanceSummary> {
111 self.instance_manager
112 .start(flow_key, context, initiator)
113 .await
114 }
115
116 pub async fn suspend_instance(&self, instance_id: &str, actor: &str) -> WorkflowResult<()> {
117 self.instance_manager.suspend(instance_id, actor).await
118 }
119
120 pub async fn resume_instance(&self, instance_id: &str, actor: &str) -> WorkflowResult<()> {
121 self.instance_manager.resume(instance_id, actor).await
122 }
123
124 pub async fn terminate_instance(&self, instance_id: &str, actor: &str) -> WorkflowResult<()> {
125 self.instance_manager.terminate(instance_id, actor).await
126 }
127
128 pub async fn query_instance(&self, instance_id: &str) -> WorkflowResult<InstanceDetail> {
129 self.instance_manager.query(instance_id).await
130 }
131
132 pub async fn query_history(&self, instance_id: &str) -> WorkflowResult<Vec<HistoryEntry>> {
133 self.instance_manager.query_history(instance_id).await
134 }
135
136 pub async fn fire_event(
137 &self,
138 instance_id: &str,
139 event: &str,
140 machine: &crate::definition::StateMachineDefinition,
141 payload: serde_json::Value,
142 ) -> WorkflowResult<TransitionResult> {
143 self.state_machine
144 .fire(instance_id, event, machine, payload)
145 .await
146 }
147
148 pub async fn handle_task(
149 &self,
150 task_id: &str,
151 action: TaskAction,
152 comment: Option<String>,
153 actor: &str,
154 definition: &crate::definition::FlowDefinition,
155 ) -> WorkflowResult<TaskHandleResult> {
156 self.approval_flow
157 .handle(task_id, action, comment, actor, definition)
158 .await
159 }
160
161 pub async fn query_tasks(
162 &self,
163 candidate: &str,
164 page: PageRequest,
165 ) -> WorkflowResult<PageResult<Task>> {
166 self.task_manager
167 .list_pending_by_candidate(candidate, page)
168 .await
169 }
170}
171
172#[cfg(test)]
173mod tests {
174 use super::*;
175 use crate::instance::InstanceStatus;
176
177 #[tokio::test]
178 async fn engine_init() {
179 let config = WorkflowConfig::default();
180 let deps = WorkflowDeps::default_for_test();
181 let engine = WorkflowEngine::new(config, deps);
182
183 let yaml = r#"
185flow_key: init_test
186version: "1.0.0"
187name: 初始化测试
188nodes:
189 - node_id: start
190 node_type: start
191 kind: start
192 next: end
193 - node_id: end
194 node_type: end
195 kind: end
196start_node: start
197active: true
198"#;
199 let result = engine.import_definition(yaml, DefinitionFormat::Yaml).await;
200 assert!(
201 result.is_ok(),
202 "engine_init 后应能导入合法流程定义,实际: {:?}",
203 result
204 );
205 }
206
207 #[tokio::test]
208 async fn full_workflow() {
209 let config = WorkflowConfig::default();
210 let deps = WorkflowDeps::default_for_test();
211 let engine = WorkflowEngine::new(config, deps);
212
213 let yaml = r#"
214flow_key: simple_flow
215version: "1.0.0"
216name: 简单流程
217nodes:
218 - node_id: start
219 node_type: start
220 kind: start
221 next: end
222 - node_id: end
223 node_type: end
224 kind: end
225start_node: start
226active: true
227"#;
228 let id = engine
229 .import_definition(yaml, DefinitionFormat::Yaml)
230 .await
231 .unwrap();
232
233 let summary = engine
234 .start_instance("simple_flow", serde_json::json!({}), "user1")
235 .await
236 .unwrap();
237 assert_eq!(summary.status, InstanceStatus::Running);
238
239 let detail = engine.query_instance(&summary.instance_id).await.unwrap();
240 assert_eq!(detail.instance.instance_id, summary.instance_id);
241
242 let exported = engine
243 .export_definition(&id, DefinitionFormat::Json)
244 .await
245 .unwrap();
246 assert!(exported.contains("simple_flow"));
247 }
248}