Skip to main content

sz_rust_workflow/scheduling/
fault_strategy.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2024-2026 SZ-Rust Team
3//
4use std::time::Duration;
5
6use crate::definition::FaultStrategy;
7use crate::error::{WorkflowError, WorkflowResult};
8
9/// 插件节点执行结果。
10#[derive(Debug, Clone, Copy, PartialEq, Eq)]
11pub enum PluginNodeOutcome {
12    Completed,
13    Skipped,
14    InstanceTerminated,
15}
16
17/// 容错策略决策。
18#[derive(Debug, Clone, Copy, PartialEq, Eq)]
19pub enum FaultDecision {
20    /// 终止实例
21    Terminate,
22    /// 跳过节点
23    Skip,
24    /// 重试(返回剩余重试次数与退避时间)
25    Retry { remaining: u32, backoff: Duration },
26}
27
28/// 容错策略处理 trait。
29pub trait FaultStrategyHandler: Send + Sync + 'static {
30    /// 根据策略与错误决定下一步动作。
31    fn decide(&self, strategy: FaultStrategy, error: &WorkflowError, attempt: u32)
32        -> FaultDecision;
33}
34
35/// 默认容错策略处理器。
36pub struct DefaultFaultStrategyHandler {
37    retry_max: u32,
38    retry_backoff: Duration,
39}
40
41impl DefaultFaultStrategyHandler {
42    pub fn new(retry_max: u32, retry_backoff: Duration) -> Self {
43        Self {
44            retry_max,
45            retry_backoff,
46        }
47    }
48}
49
50impl Default for DefaultFaultStrategyHandler {
51    fn default() -> Self {
52        Self::new(3, Duration::from_millis(100))
53    }
54}
55
56impl FaultStrategyHandler for DefaultFaultStrategyHandler {
57    fn decide(
58        &self,
59        strategy: FaultStrategy,
60        error: &WorkflowError,
61        attempt: u32,
62    ) -> FaultDecision {
63        match strategy {
64            FaultStrategy::Fail => {
65                tracing::error!(error = %error, "插件节点失败,终止实例");
66                FaultDecision::Terminate
67            }
68            FaultStrategy::Skip => {
69                tracing::warn!(error = %error, "插件节点失败,跳过");
70                FaultDecision::Skip
71            }
72            FaultStrategy::Retry => {
73                if attempt < self.retry_max {
74                    let backoff = self.retry_backoff * 2u32.pow(attempt);
75                    tracing::warn!(
76                        attempt = attempt + 1,
77                        backoff_ms = backoff.as_millis(),
78                        "重试"
79                    );
80                    FaultDecision::Retry {
81                        remaining: self.retry_max - attempt,
82                        backoff,
83                    }
84                } else {
85                    tracing::error!(error = %error, "重试超限,终止实例");
86                    FaultDecision::Terminate
87                }
88            }
89        }
90    }
91}
92
93/// 便捷函数:根据决策返回 PluginNodeOutcome。
94pub fn outcome_from_decision(decision: FaultDecision) -> WorkflowResult<PluginNodeOutcome> {
95    match decision {
96        FaultDecision::Terminate => Ok(PluginNodeOutcome::InstanceTerminated),
97        FaultDecision::Skip => Ok(PluginNodeOutcome::Skipped),
98        FaultDecision::Retry { remaining: 0, .. } => Ok(PluginNodeOutcome::InstanceTerminated),
99        FaultDecision::Retry { .. } => Ok(PluginNodeOutcome::Completed),
100    }
101}
102
103#[cfg(test)]
104mod tests {
105    use super::*;
106    use crate::error::WorkflowErrorCode;
107
108    #[test]
109    fn fail_strategy() {
110        let handler = DefaultFaultStrategyHandler::default();
111        let error = WorkflowError::new(WorkflowErrorCode::CapabilityNotFound, "能力不存在");
112        let decision = handler.decide(FaultStrategy::Fail, &error, 0);
113        assert_eq!(decision, FaultDecision::Terminate);
114    }
115
116    #[test]
117    fn skip_strategy() {
118        let handler = DefaultFaultStrategyHandler::default();
119        let error = WorkflowError::new(WorkflowErrorCode::CapabilityNotFound, "能力不存在");
120        let decision = handler.decide(FaultStrategy::Skip, &error, 0);
121        assert_eq!(decision, FaultDecision::Skip);
122    }
123
124    #[test]
125    fn retry_first_attempt() {
126        let handler = DefaultFaultStrategyHandler::new(3, Duration::from_millis(10));
127        let error = WorkflowError::new(WorkflowErrorCode::CapabilityNotFound, "能力不存在");
128        let decision = handler.decide(FaultStrategy::Retry, &error, 0);
129        match decision {
130            FaultDecision::Retry { remaining, backoff } => {
131                assert_eq!(remaining, 3);
132                assert_eq!(backoff, Duration::from_millis(10));
133            }
134            _ => panic!("期望 Retry"),
135        }
136    }
137
138    #[test]
139    fn retry_second_attempt() {
140        let handler = DefaultFaultStrategyHandler::new(3, Duration::from_millis(10));
141        let error = WorkflowError::new(WorkflowErrorCode::CapabilityNotFound, "能力不存在");
142        let decision = handler.decide(FaultStrategy::Retry, &error, 1);
143        match decision {
144            FaultDecision::Retry { remaining, backoff } => {
145                assert_eq!(remaining, 2);
146                assert_eq!(backoff, Duration::from_millis(20));
147            }
148            _ => panic!("期望 Retry"),
149        }
150    }
151
152    #[test]
153    fn retry_exhausted() {
154        let handler = DefaultFaultStrategyHandler::new(2, Duration::from_millis(1));
155        let error = WorkflowError::new(WorkflowErrorCode::CapabilityNotFound, "能力不存在");
156        let decision = handler.decide(FaultStrategy::Retry, &error, 2);
157        assert_eq!(decision, FaultDecision::Terminate);
158    }
159
160    #[test]
161    fn outcome_terminate() {
162        let result = outcome_from_decision(FaultDecision::Terminate);
163        assert_eq!(result.unwrap(), PluginNodeOutcome::InstanceTerminated);
164    }
165
166    #[test]
167    fn outcome_skip() {
168        let result = outcome_from_decision(FaultDecision::Skip);
169        assert_eq!(result.unwrap(), PluginNodeOutcome::Skipped);
170    }
171
172    #[test]
173    fn outcome_retry_with_remaining() {
174        let decision = FaultDecision::Retry {
175            remaining: 2,
176            backoff: Duration::from_millis(10),
177        };
178        let result = outcome_from_decision(decision);
179        assert_eq!(result.unwrap(), PluginNodeOutcome::Completed);
180    }
181
182    #[test]
183    fn outcome_retry_zero_remaining() {
184        let decision = FaultDecision::Retry {
185            remaining: 0,
186            backoff: Duration::from_millis(10),
187        };
188        let result = outcome_from_decision(decision);
189        assert_eq!(result.unwrap(), PluginNodeOutcome::InstanceTerminated);
190    }
191}