1use crate::WorkGraphError;
2use crate::generated::{
3 protocol_work_execution_cancellation_evidence_projection as cancellation_evidence_protocol,
4 protocol_work_execution_failure_evidence_projection as failure_evidence_protocol,
5 protocol_work_execution_flow_launch as flow_launch_protocol,
6 protocol_work_execution_flow_observation as flow_observation_protocol,
7 protocol_work_execution_launch_failure_evidence_projection as launch_failure_evidence_protocol,
8 protocol_work_execution_quarantined_launch_resolution as quarantined_launch_protocol,
9 protocol_work_execution_success_evidence_projection as success_evidence_protocol,
10 protocol_work_execution_uncertain_launch_resolution as uncertain_launch_protocol,
11 protocol_work_execution_work_closure as work_closure_protocol,
12};
13use crate::machines::work_execution_lifecycle as execution_dsl;
14use crate::types::{WorkExecutionBinding, WorkExecutionBindingId, WorkExecutionMachineState};
15
16pub use execution_dsl::{WorkExecutionLifecycleEffect, WorkExecutionLifecycleState};
17
18#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
19pub enum WorkExecutionObservation {
20 FlowStarted,
21 FlowRunning,
22 FlowCompleted,
23 FlowFailed { detail: Option<String> },
24 FlowCanceled { detail: Option<String> },
25 FlowRunLost { detail: String },
26 LaunchUncertain { detail: String },
27 LaunchQuarantined { detail: String },
28 LaunchFailed { detail: String },
29 EvidenceProjected,
30 FlowFailureEvidenceProjected,
31 FlowCancellationEvidenceProjected,
32 LaunchFailureEvidenceProjected,
33 WorkClosed,
34 WorkClosureRefused { detail: String },
35}
36
37#[derive(Debug, Default, Clone, Copy)]
38pub struct WorkExecutionMachine;
39
40#[derive(Debug, Clone, PartialEq, Eq)]
41pub struct WorkExecutionTransition {
42 pub binding: WorkExecutionBinding,
43 pub effect: WorkExecutionLifecycleEffect,
44}
45
46pub struct WorkExecutionBindCommit {
48 binding: WorkExecutionBinding,
49 effect: WorkExecutionLifecycleEffect,
50}
51
52impl WorkExecutionBindCommit {
53 pub(crate) fn into_parts(self) -> (WorkExecutionBinding, WorkExecutionLifecycleEffect) {
54 (self.binding, self.effect)
55 }
56
57 pub(crate) fn binding(&self) -> &WorkExecutionBinding {
58 &self.binding
59 }
60
61 pub(crate) fn effect(&self) -> &WorkExecutionLifecycleEffect {
62 &self.effect
63 }
64}
65
66pub struct WorkExecutionObservationCommit {
68 previous: WorkExecutionBinding,
69 observation: WorkExecutionObservation,
70 binding: WorkExecutionBinding,
71 effect: WorkExecutionLifecycleEffect,
72}
73
74impl WorkExecutionObservationCommit {
75 pub(crate) fn into_parts(
76 self,
77 ) -> (
78 WorkExecutionBinding,
79 WorkExecutionObservation,
80 WorkExecutionBinding,
81 WorkExecutionLifecycleEffect,
82 ) {
83 (self.previous, self.observation, self.binding, self.effect)
84 }
85
86 pub(crate) fn binding(&self) -> &WorkExecutionBinding {
87 &self.binding
88 }
89
90 pub(crate) fn effect(&self) -> &WorkExecutionLifecycleEffect {
91 &self.effect
92 }
93}
94
95impl WorkExecutionMachine {
96 pub(crate) fn prepare_bind(
97 binding: WorkExecutionBinding,
98 ) -> Result<WorkExecutionBindCommit, WorkGraphError> {
99 binding.validate()?;
100 let (expected, effect) = Self::bind(&binding.binding_id, binding.target.run_id())?;
101 if binding.machine_state != expected {
102 return Err(WorkGraphError::InvalidInput(format!(
103 "work execution binding {} was not initialized by WorkExecutionLifecycleMachine",
104 binding.binding_id
105 )));
106 }
107 Ok(WorkExecutionBindCommit { binding, effect })
108 }
109
110 pub(crate) fn prepare_observation(
111 previous: WorkExecutionBinding,
112 expected_revision: u64,
113 observation: WorkExecutionObservation,
114 ) -> Result<WorkExecutionObservationCommit, WorkGraphError> {
115 let (binding, effect) =
116 Self::observe(previous.clone(), expected_revision, observation.clone())?;
117 Ok(WorkExecutionObservationCommit {
118 previous,
119 observation,
120 binding,
121 effect,
122 })
123 }
124 pub fn recover_effect(
125 binding: &WorkExecutionBinding,
126 ) -> Result<WorkExecutionLifecycleEffect, WorkGraphError> {
127 validate_projection(binding)?;
128 let mut authority =
129 execution_dsl::WorkExecutionLifecycleMachineAuthority::recover_from_state(
130 binding.machine_state.clone(),
131 )
132 .map_err(|error| {
133 WorkGraphError::InvalidTransition(format!(
134 "work execution {} refused recovery: {error:?}",
135 binding.binding_id
136 ))
137 })?;
138 let transition = execution_dsl::WorkExecutionLifecycleMachineMutator::apply(
139 &mut authority,
140 execution_dsl::WorkExecutionLifecycleInput::Recover {},
141 )
142 .map_err(|error| {
143 WorkGraphError::InvalidTransition(format!(
144 "work execution {} refused effect recovery: {error:?}",
145 binding.binding_id
146 ))
147 })?;
148 validate_handoff_obligation(&transition)?;
149 exactly_one_effect(transition.effects())
150 }
151
152 pub fn bind(
153 binding_id: &WorkExecutionBindingId,
154 run_id: &str,
155 ) -> Result<(WorkExecutionMachineState, WorkExecutionLifecycleEffect), WorkGraphError> {
156 let mut authority = execution_dsl::WorkExecutionLifecycleMachineAuthority::new();
157 let transition = execution_dsl::WorkExecutionLifecycleMachineMutator::apply(
158 &mut authority,
159 execution_dsl::WorkExecutionLifecycleInput::Bind {
160 binding_id: binding_id.as_str().to_string(),
161 run_id: run_id.to_string(),
162 },
163 )
164 .map_err(|error| {
165 WorkGraphError::InvalidTransition(format!(
166 "generated work execution binding transition refused: {error:?}"
167 ))
168 })?;
169 validate_handoff_obligation(&transition)?;
170 let effect = exactly_one_effect(transition.effects())?;
171 Ok((authority.state().clone(), effect))
172 }
173
174 pub fn observe(
175 mut binding: WorkExecutionBinding,
176 expected_revision: u64,
177 observation: WorkExecutionObservation,
178 ) -> Result<(WorkExecutionBinding, WorkExecutionLifecycleEffect), WorkGraphError> {
179 validate_observation_detail(&observation)?;
180 if binding.machine_state.revision != expected_revision {
181 return Err(WorkGraphError::Conflict(format!(
182 "stale work execution revision for {}: expected {}, actual {}",
183 binding.binding_id, expected_revision, binding.machine_state.revision
184 )));
185 }
186 validate_projection(&binding)?;
187 let mut authority =
188 execution_dsl::WorkExecutionLifecycleMachineAuthority::recover_from_state(
189 binding.machine_state.clone(),
190 )
191 .map_err(|error| {
192 WorkGraphError::InvalidTransition(format!(
193 "work execution {} refused recovery: {error:?}",
194 binding.binding_id
195 ))
196 })?;
197 let mut recovery_authority =
198 execution_dsl::WorkExecutionLifecycleMachineAuthority::recover_from_state(
199 binding.machine_state.clone(),
200 )
201 .map_err(|error| {
202 WorkGraphError::InvalidTransition(format!(
203 "work execution {} refused handoff recovery: {error:?}",
204 binding.binding_id
205 ))
206 })?;
207 let recovery_transition = execution_dsl::WorkExecutionLifecycleMachineMutator::apply(
208 &mut recovery_authority,
209 execution_dsl::WorkExecutionLifecycleInput::Recover {},
210 )
211 .map_err(|error| {
212 WorkGraphError::InvalidTransition(format!(
213 "work execution {} refused handoff effect recovery: {error:?}",
214 binding.binding_id
215 ))
216 })?;
217 validate_handoff_obligation(&recovery_transition)?;
218 let recovered_effect = exactly_one_effect(recovery_transition.effects())?;
219 let observation_debug = format!("{observation:?}");
220 macro_rules! submit {
221 ($protocol:ident, $function:ident $(, $argument:expr)*) => {{
222 let obligation = exactly_one_obligation(
223 $protocol::extract_obligations(&recovery_transition),
224 &binding.binding_id,
225 )?;
226 $protocol::$function(&mut authority, obligation $(, $argument)*)
227 }};
228 }
229 let transition = match (&recovered_effect, observation) {
230 (WorkExecutionLifecycleEffect::FlowLaunchRequested { .. }, WorkExecutionObservation::FlowStarted) =>
231 submit!(flow_launch_protocol, submit_confirm_flow_started),
232 (WorkExecutionLifecycleEffect::FlowLaunchRequested { .. }, WorkExecutionObservation::FlowRunning) =>
233 submit!(flow_launch_protocol, submit_observe_flow_running),
234 (WorkExecutionLifecycleEffect::FlowLaunchRequested { .. }, WorkExecutionObservation::FlowCompleted) =>
235 submit!(flow_launch_protocol, submit_observe_flow_completed),
236 (WorkExecutionLifecycleEffect::FlowLaunchRequested { .. }, WorkExecutionObservation::FlowFailed { detail }) =>
237 submit!(flow_launch_protocol, submit_observe_flow_failed, detail),
238 (WorkExecutionLifecycleEffect::FlowLaunchRequested { .. }, WorkExecutionObservation::FlowCanceled { detail }) =>
239 submit!(flow_launch_protocol, submit_observe_flow_canceled, detail),
240 (WorkExecutionLifecycleEffect::FlowLaunchRequested { .. }, WorkExecutionObservation::LaunchUncertain { detail }) =>
241 submit!(flow_launch_protocol, submit_mark_launch_uncertain, detail),
242 (WorkExecutionLifecycleEffect::FlowLaunchRequested { .. }, WorkExecutionObservation::LaunchQuarantined { detail }) =>
243 submit!(flow_launch_protocol, submit_quarantine_launch, detail),
244 (WorkExecutionLifecycleEffect::FlowLaunchRequested { .. }, WorkExecutionObservation::LaunchFailed { detail }) =>
245 submit!(flow_launch_protocol, submit_resolve_launch_failed, detail),
246
247 (WorkExecutionLifecycleEffect::FlowLaunchAccepted { .. }, WorkExecutionObservation::FlowRunning) =>
248 submit!(flow_observation_protocol, submit_observe_flow_running),
249 (WorkExecutionLifecycleEffect::FlowLaunchAccepted { .. }, WorkExecutionObservation::FlowCompleted) =>
250 submit!(flow_observation_protocol, submit_observe_flow_completed),
251 (WorkExecutionLifecycleEffect::FlowLaunchAccepted { .. }, WorkExecutionObservation::FlowFailed { detail }) =>
252 submit!(flow_observation_protocol, submit_observe_flow_failed, detail),
253 (WorkExecutionLifecycleEffect::FlowLaunchAccepted { .. }, WorkExecutionObservation::FlowCanceled { detail }) =>
254 submit!(flow_observation_protocol, submit_observe_flow_canceled, detail),
255 (WorkExecutionLifecycleEffect::FlowLaunchAccepted { .. }, WorkExecutionObservation::FlowRunLost { detail }) =>
256 submit!(flow_observation_protocol, submit_observe_run_lost, detail),
257
258 (WorkExecutionLifecycleEffect::FlowLaunchUncertain { .. }, WorkExecutionObservation::FlowStarted) =>
259 submit!(uncertain_launch_protocol, submit_confirm_flow_started),
260 (WorkExecutionLifecycleEffect::FlowLaunchUncertain { .. }, WorkExecutionObservation::FlowRunning) =>
261 submit!(uncertain_launch_protocol, submit_observe_flow_running),
262 (WorkExecutionLifecycleEffect::FlowLaunchUncertain { .. }, WorkExecutionObservation::FlowCompleted) =>
263 submit!(uncertain_launch_protocol, submit_observe_flow_completed),
264 (WorkExecutionLifecycleEffect::FlowLaunchUncertain { .. }, WorkExecutionObservation::FlowFailed { detail }) =>
265 submit!(uncertain_launch_protocol, submit_observe_flow_failed, detail),
266 (WorkExecutionLifecycleEffect::FlowLaunchUncertain { .. }, WorkExecutionObservation::FlowCanceled { detail }) =>
267 submit!(uncertain_launch_protocol, submit_observe_flow_canceled, detail),
268 (WorkExecutionLifecycleEffect::FlowLaunchUncertain { .. }, WorkExecutionObservation::LaunchFailed { detail }) =>
269 submit!(uncertain_launch_protocol, submit_resolve_launch_failed, detail),
270 (WorkExecutionLifecycleEffect::FlowLaunchUncertain { .. }, WorkExecutionObservation::LaunchQuarantined { detail }) =>
271 submit!(uncertain_launch_protocol, submit_quarantine_launch, detail),
272
273 (WorkExecutionLifecycleEffect::FlowLaunchQuarantined { .. }, WorkExecutionObservation::FlowStarted) =>
274 submit!(quarantined_launch_protocol, submit_confirm_flow_started),
275 (WorkExecutionLifecycleEffect::FlowLaunchQuarantined { .. }, WorkExecutionObservation::FlowRunning) =>
276 submit!(quarantined_launch_protocol, submit_observe_flow_running),
277 (WorkExecutionLifecycleEffect::FlowLaunchQuarantined { .. }, WorkExecutionObservation::FlowCompleted) =>
278 submit!(quarantined_launch_protocol, submit_observe_flow_completed),
279 (WorkExecutionLifecycleEffect::FlowLaunchQuarantined { .. }, WorkExecutionObservation::FlowFailed { detail }) =>
280 submit!(quarantined_launch_protocol, submit_observe_flow_failed, detail),
281 (WorkExecutionLifecycleEffect::FlowLaunchQuarantined { .. }, WorkExecutionObservation::FlowCanceled { detail }) =>
282 submit!(quarantined_launch_protocol, submit_observe_flow_canceled, detail),
283
284 (WorkExecutionLifecycleEffect::EvidenceProjectionRequested { .. }, WorkExecutionObservation::EvidenceProjected) =>
285 submit!(success_evidence_protocol, submit_confirm_evidence_projected),
286 (WorkExecutionLifecycleEffect::EvidenceProjectionRequested { .. }, WorkExecutionObservation::FlowRunLost { detail }) =>
287 submit!(success_evidence_protocol, submit_observe_run_lost, detail),
288 (WorkExecutionLifecycleEffect::FlowFailureEvidenceProjectionRequested { .. }, WorkExecutionObservation::FlowFailureEvidenceProjected) =>
289 submit!(failure_evidence_protocol, submit_confirm_flow_failure_evidence_projected),
290 (WorkExecutionLifecycleEffect::FlowCancellationEvidenceProjectionRequested { .. }, WorkExecutionObservation::FlowCancellationEvidenceProjected) =>
291 submit!(cancellation_evidence_protocol, submit_confirm_flow_cancellation_evidence_projected),
292 (WorkExecutionLifecycleEffect::LaunchFailureEvidenceProjectionRequested { .. }, WorkExecutionObservation::LaunchFailureEvidenceProjected) =>
293 submit!(launch_failure_evidence_protocol, submit_confirm_launch_failure_evidence_projected),
294 (WorkExecutionLifecycleEffect::WorkClosureRequested { .. }, WorkExecutionObservation::WorkClosed) =>
295 submit!(work_closure_protocol, submit_confirm_work_closed),
296 (WorkExecutionLifecycleEffect::WorkClosureRequested { .. }, WorkExecutionObservation::WorkClosureRefused { detail }) =>
297 submit!(work_closure_protocol, submit_refuse_work_closure, detail),
298 _ => {
299 return Err(WorkGraphError::InvalidTransition(format!(
300 "work execution {} observation {observation_debug} is not admitted by the generated owner-feedback protocol for {recovered_effect:?}",
301 binding.binding_id
302 )));
303 }
304 }
305 .map_err(|error| {
306 WorkGraphError::InvalidTransition(format!(
307 "work execution {} refused observation: {error:?}",
308 binding.binding_id
309 ))
310 })?;
311 validate_handoff_obligation(&transition)?;
312 let effect = exactly_one_effect(transition.effects())?;
313 binding.machine_state = authority.state().clone();
314 validate_projection(&binding)?;
315 Ok((binding, effect))
316 }
317
318 pub fn validate_projection(binding: &WorkExecutionBinding) -> Result<(), WorkGraphError> {
319 validate_projection(binding)
320 }
321
322 pub fn retry_eligible(binding: &WorkExecutionBinding) -> Result<bool, WorkGraphError> {
324 validate_projection(binding)?;
325 let mut authority =
326 execution_dsl::WorkExecutionLifecycleMachineAuthority::recover_from_state(
327 binding.machine_state.clone(),
328 )
329 .map_err(|error| {
330 WorkGraphError::InvalidTransition(format!(
331 "work execution {} refused retry classification recovery: {error:?}",
332 binding.binding_id
333 ))
334 })?;
335 let transition = execution_dsl::WorkExecutionLifecycleMachineMutator::apply(
336 &mut authority,
337 execution_dsl::WorkExecutionLifecycleInput::ClassifyRetryEligibility {},
338 )
339 .map_err(|error| {
340 WorkGraphError::InvalidTransition(format!(
341 "work execution {} refused retry classification: {error:?}",
342 binding.binding_id
343 ))
344 })?;
345 match exactly_one_effect(transition.effects())? {
346 WorkExecutionLifecycleEffect::RetryEligibilityClassified { eligible } => Ok(eligible),
347 effect => Err(WorkGraphError::Store(format!(
348 "work execution {} retry classifier emitted unexpected effect {effect:?}",
349 binding.binding_id
350 ))),
351 }
352 }
353}
354
355fn validate_observation_detail(
356 observation: &WorkExecutionObservation,
357) -> Result<(), WorkGraphError> {
358 const MAX_DETAIL_BYTES: usize = 4096;
359 let detail = match observation {
360 WorkExecutionObservation::FlowFailed { detail }
361 | WorkExecutionObservation::FlowCanceled { detail } => detail.as_deref(),
362 WorkExecutionObservation::FlowRunLost { detail }
363 | WorkExecutionObservation::LaunchUncertain { detail }
364 | WorkExecutionObservation::LaunchQuarantined { detail }
365 | WorkExecutionObservation::LaunchFailed { detail }
366 | WorkExecutionObservation::WorkClosureRefused { detail } => Some(detail.as_str()),
367 _ => None,
368 };
369 if let Some(detail) = detail
370 && (detail.len() > MAX_DETAIL_BYTES || detail.chars().any(char::is_control))
371 {
372 return Err(WorkGraphError::InvalidInput(format!(
373 "work execution observation detail must be single-line text no longer than {MAX_DETAIL_BYTES} bytes"
374 )));
375 }
376 Ok(())
377}
378
379fn validate_handoff_obligation(
380 transition: &execution_dsl::WorkExecutionLifecycleMachineTransition,
381) -> Result<(), WorkGraphError> {
382 let obligation_count = match transition.effects() {
383 [WorkExecutionLifecycleEffect::FlowLaunchRequested { .. }] => {
384 flow_launch_protocol::extract_obligations(transition).len()
385 }
386 [WorkExecutionLifecycleEffect::FlowLaunchAccepted { .. }] => {
387 flow_observation_protocol::extract_obligations(transition).len()
388 }
389 [WorkExecutionLifecycleEffect::FlowLaunchUncertain { .. }] => {
390 uncertain_launch_protocol::extract_obligations(transition).len()
391 }
392 [WorkExecutionLifecycleEffect::FlowLaunchQuarantined { .. }] => {
393 quarantined_launch_protocol::extract_obligations(transition).len()
394 }
395 [WorkExecutionLifecycleEffect::EvidenceProjectionRequested { .. }] => {
396 success_evidence_protocol::extract_obligations(transition).len()
397 }
398 [WorkExecutionLifecycleEffect::FlowFailureEvidenceProjectionRequested { .. }] => {
399 failure_evidence_protocol::extract_obligations(transition).len()
400 }
401 [WorkExecutionLifecycleEffect::FlowCancellationEvidenceProjectionRequested { .. }] => {
402 cancellation_evidence_protocol::extract_obligations(transition).len()
403 }
404 [WorkExecutionLifecycleEffect::LaunchFailureEvidenceProjectionRequested { .. }] => {
405 launch_failure_evidence_protocol::extract_obligations(transition).len()
406 }
407 [WorkExecutionLifecycleEffect::WorkClosureRequested { .. }] => {
408 work_closure_protocol::extract_obligations(transition).len()
409 }
410 _ => return Ok(()),
411 };
412 if obligation_count != 1 {
413 return Err(WorkGraphError::Store(format!(
414 "generated work execution transition projected {obligation_count} handoff obligations, expected exactly one"
415 )));
416 }
417 Ok(())
418}
419
420fn validate_projection(binding: &WorkExecutionBinding) -> Result<(), WorkGraphError> {
421 if binding.machine_state.binding_id != binding.binding_id.as_str() {
422 return Err(WorkGraphError::Store(format!(
423 "work execution {} machine binding identity does not match its projection",
424 binding.binding_id
425 )));
426 }
427 if binding.machine_state.run_id != binding.target.run_id() {
428 return Err(WorkGraphError::Store(format!(
429 "work execution {} machine run identity does not match its target",
430 binding.binding_id
431 )));
432 }
433 Ok(())
434}
435
436fn exactly_one_effect(
437 effects: &[WorkExecutionLifecycleEffect],
438) -> Result<WorkExecutionLifecycleEffect, WorkGraphError> {
439 match effects {
440 [effect] => Ok(effect.clone()),
441 _ => Err(WorkGraphError::Store(format!(
442 "generated work execution transition emitted {} effects, expected exactly one",
443 effects.len()
444 ))),
445 }
446}
447
448fn exactly_one_obligation<T>(
449 obligations: Vec<T>,
450 binding_id: &WorkExecutionBindingId,
451) -> Result<T, WorkGraphError> {
452 let count = obligations.len();
453 obligations.into_iter().next().filter(|_| count == 1).ok_or_else(|| {
454 WorkGraphError::Store(format!(
455 "generated work execution handoff for {binding_id} projected {count} obligations, expected exactly one"
456 ))
457 })
458}
459
460#[cfg(test)]
461#[allow(clippy::expect_used)]
462mod tests {
463 use super::*;
464 use crate::{WorkExecutionTarget, WorkItemId, WorkItemRef, WorkNamespace};
465 use chrono::Utc;
466 use serde_json::json;
467
468 fn binding() -> WorkExecutionBinding {
469 let binding_id = WorkExecutionBindingId::new("execution-machine-test").expect("id");
470 let target = WorkExecutionTarget::mob_flow(
471 "mob",
472 "flow",
473 format!("sha256:{}", "c".repeat(64)),
474 "1ae92ab4-8afe-5ad2-b9c3-fccae4f569a5",
475 crate::WorkExecutionAuthority::TargetOwner,
476 json!({}),
477 )
478 .expect("target");
479 let (machine_state, _) =
480 WorkExecutionMachine::bind(&binding_id, target.run_id()).expect("bind");
481 WorkExecutionBinding {
482 binding_id,
483 work_ref: WorkItemRef {
484 realm_id: "realm".to_string(),
485 namespace: WorkNamespace::default(),
486 item_id: WorkItemId::new("item").expect("item"),
487 },
488 target,
489 idempotency_key: "key".to_string(),
490 correlation_id: "229650c5-9372-53e9-9c3a-831638a47c77".to_string(),
491 supersedes: None,
492 machine_state,
493 created_at: Utc::now(),
494 }
495 }
496
497 #[test]
498 fn failed_flow_requires_evidence_projection_before_terminal_attempt() {
499 let binding = binding();
500 let (running, _) =
501 WorkExecutionMachine::observe(binding, 1, WorkExecutionObservation::FlowRunning)
502 .expect("running");
503 let (projecting, effect) = WorkExecutionMachine::observe(
504 running,
505 2,
506 WorkExecutionObservation::FlowFailed {
507 detail: Some("step failed".to_string()),
508 },
509 )
510 .expect("failure observed");
511 assert!(matches!(
512 effect,
513 WorkExecutionLifecycleEffect::FlowFailureEvidenceProjectionRequested { .. }
514 ));
515 let (terminal, effect) = WorkExecutionMachine::observe(
516 projecting,
517 3,
518 WorkExecutionObservation::FlowFailureEvidenceProjected,
519 )
520 .expect("failure evidence projected");
521 assert!(matches!(
522 effect,
523 WorkExecutionLifecycleEffect::FlowFailed { .. }
524 ));
525 assert!(matches!(
526 WorkExecutionMachine::recover_effect(&terminal).expect("recover terminal"),
527 WorkExecutionLifecycleEffect::FlowFailed { .. }
528 ));
529 }
530
531 #[test]
532 fn uncertain_launch_recovers_as_uncertain_without_redrive_request() {
533 let binding = binding();
534 let (uncertain, effect) = WorkExecutionMachine::observe(
535 binding,
536 1,
537 WorkExecutionObservation::LaunchUncertain {
538 detail: "intent exists without run".to_string(),
539 },
540 )
541 .expect("uncertain");
542 assert!(matches!(
543 effect,
544 WorkExecutionLifecycleEffect::FlowLaunchUncertain { .. }
545 ));
546 assert!(matches!(
547 WorkExecutionMachine::recover_effect(&uncertain).expect("recover uncertain"),
548 WorkExecutionLifecycleEffect::FlowLaunchUncertain { .. }
549 ));
550 }
551
552 #[test]
553 fn observations_are_admitted_only_by_the_current_generated_handoff() {
554 let binding = binding();
555 let (accepted, _) =
556 WorkExecutionMachine::observe(binding, 1, WorkExecutionObservation::FlowStarted)
557 .expect("launch accepted");
558 let error =
559 WorkExecutionMachine::observe(accepted, 2, WorkExecutionObservation::FlowStarted)
560 .expect_err("flow-observation protocol does not admit a second start");
561 assert!(matches!(error, WorkGraphError::InvalidTransition(_)));
562 }
563
564 #[test]
565 fn quarantined_launch_accepts_only_observed_exact_run_feedback() {
566 let binding = binding();
567 let (quarantined, effect) = WorkExecutionMachine::observe(
568 binding,
569 1,
570 WorkExecutionObservation::LaunchQuarantined {
571 detail: "realizing ledger has no exact run".to_string(),
572 },
573 )
574 .expect("quarantine");
575 assert!(matches!(
576 effect,
577 WorkExecutionLifecycleEffect::FlowLaunchQuarantined { .. }
578 ));
579 assert!(
580 !WorkExecutionMachine::retry_eligible(&quarantined)
581 .expect("classify quarantined launch"),
582 "quarantine may still have a live realizer and cannot authorize supersession"
583 );
584 let (_, effect) =
585 WorkExecutionMachine::observe(quarantined, 2, WorkExecutionObservation::FlowStarted)
586 .expect("exact run observation feedback");
587 assert!(matches!(
588 effect,
589 WorkExecutionLifecycleEffect::FlowLaunchAccepted { .. }
590 ));
591 }
592}