Skip to main content

meerkat_workgraph/
execution_machine.rs

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
46/// Sealed machine-minted authority for the initial durable binding commit.
47pub 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
66/// Sealed machine-minted authority for one durable observation transition.
67pub 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    /// The lifecycle authority's single retry/supersession classifier.
323    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}