sz_rust_workflow/scheduling/
fault_strategy.rs1use std::time::Duration;
5
6use crate::definition::FaultStrategy;
7use crate::error::{WorkflowError, WorkflowResult};
8
9#[derive(Debug, Clone, Copy, PartialEq, Eq)]
11pub enum PluginNodeOutcome {
12 Completed,
13 Skipped,
14 InstanceTerminated,
15}
16
17#[derive(Debug, Clone, Copy, PartialEq, Eq)]
19pub enum FaultDecision {
20 Terminate,
22 Skip,
24 Retry { remaining: u32, backoff: Duration },
26}
27
28pub trait FaultStrategyHandler: Send + Sync + 'static {
30 fn decide(&self, strategy: FaultStrategy, error: &WorkflowError, attempt: u32)
32 -> FaultDecision;
33}
34
35pub 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
93pub 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}