Skip to main content

sz_rust_workflow/engine/
workflow_engine.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2024-2026 SZ-Rust Team
3//
4use 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
19/// 工作流引擎统一门面,对齐 design 2.2.2.1。
20pub 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)] // 构造时移交给 instance_manager,字段仅持有引用
29    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        // 验证引擎可导入并校验一个合法流程定义(构造后基础能力可用)
184        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}