Skip to main content

af_workflow/
spec_driver.rs

1//! Bridge from a validated product `Spec` to the durable supervisor.
2//!
3//! Product specs get durable execution without a host-side state machine.
4//! Action steps emit data-only requests; this driver pins and persists them as
5//! `ActionIntent`s before the product provider can run.
6
7use std::collections::BTreeMap;
8use std::sync::Arc;
9
10use async_trait::async_trait;
11use serde_json::{json, Map, Value};
12
13use crate::{
14    ingress_plans, schedule_decision, ActionIntent, ActionState, CapabilityKind,
15    CapabilityManifest, CatchUpPolicy, EvaluationOutcome, Event, HostError, MemoryState,
16    NodeRegistry, ScheduleCadence, SchedulePolicy, Spec, Terminal, WorkDisposition, WorkItem,
17    WorkflowDriver, WorkflowHost, WorkflowTransitionCommand,
18};
19
20/// Seconds an event-driven instance sleeps between wakeups when no timer or
21/// trigger delivery is pending. Deliveries pull `next_run_at` forward, so this
22/// only bounds how long a stale row stays idle.
23const IDLE_WAIT_SECS: i64 = 86_400;
24
25/// Durable driver compiled from a product spec; see the module docs.
26pub struct SpecDriver {
27    spec_id: String,
28    host: WorkflowHost,
29    schedules: BTreeMap<String, SchedulePolicy>,
30    actions: BTreeMap<String, CapabilityManifest>,
31}
32
33#[derive(serde::Serialize, serde::Deserialize)]
34struct ScheduleCursor {
35    next_at: chrono::DateTime<chrono::Utc>,
36    catch_up_count: u32,
37}
38
39#[derive(serde::Serialize, serde::Deserialize)]
40struct ActionContinuation {
41    branch_id: String,
42    step: usize,
43    events: Vec<Event>,
44    #[serde(default)]
45    action_ids: Vec<String>,
46}
47
48impl SpecDriver {
49    /// Compile `spec`; rejects action nodes without an `Action` manifest and funds actions.
50    pub fn new(spec: &Spec, registry: &NodeRegistry) -> Result<Self, HostError> {
51        let host = WorkflowHost::from_spec(spec, registry)?;
52        let schedules = ingress_plans(spec, registry)
53            .into_iter()
54            .map(|plan| {
55                let branch_id = plan.branch_id.clone();
56                (|| {
57                    let required_u64 = |name: &str| {
58                        plan.ingress_config[name]
59                            .as_u64()
60                            .ok_or_else(|| HostError::Schedule {
61                                spec_id: spec.spec_id.clone(),
62                                reason: format!("{} requires integer '{name}'", plan.ingress_type),
63                            })
64                    };
65                    let cadence = match plan.ingress_type.as_str() {
66                        "ingress.cron" => ScheduleCadence::Cron {
67                            expression: plan.ingress_config["expression"]
68                                .as_str()
69                                .ok_or_else(|| HostError::Schedule {
70                                    spec_id: spec.spec_id.clone(),
71                                    reason: "ingress.cron requires a string 'expression'".into(),
72                                })?
73                                .to_owned(),
74                        },
75                        "ingress.fixed_rate" => ScheduleCadence::FixedRate {
76                            milliseconds: required_u64("milliseconds")?,
77                        },
78                        "ingress.fixed_delay" => ScheduleCadence::FixedDelay {
79                            milliseconds: required_u64("milliseconds")?,
80                        },
81                        _ => return Ok(None),
82                    };
83                    let catch_up = match plan.ingress_config["catch_up"].as_str().unwrap_or("once")
84                    {
85                        "skip" => CatchUpPolicy::Skip,
86                        "once" => CatchUpPolicy::CatchUpOnce,
87                        "all" => CatchUpPolicy::CatchUpAll {
88                            limit: plan.ingress_config["catch_up_limit"]
89                                .as_u64()
90                                .and_then(|value| u32::try_from(value).ok())
91                                .ok_or_else(|| HostError::Schedule {
92                                    spec_id: spec.spec_id.clone(),
93                                    reason: "catch_up=all requires positive integer catch_up_limit"
94                                        .into(),
95                                })?,
96                        },
97                        value => {
98                            return Err(HostError::Schedule {
99                                spec_id: spec.spec_id.clone(),
100                                reason: format!("unknown catch_up policy '{value}'"),
101                            })
102                        }
103                    };
104                    let policy = SchedulePolicy {
105                        cadence,
106                        timezone: plan.ingress_config["timezone"]
107                            .as_str()
108                            .unwrap_or("UTC")
109                            .to_owned(),
110                        catch_up,
111                    };
112                    policy.validate().map_err(|error| HostError::Schedule {
113                        spec_id: spec.spec_id.clone(),
114                        reason: error.to_string(),
115                    })?;
116                    Ok::<_, HostError>(Some(policy))
117                })()
118                .map(|policy| policy.map(|policy| (branch_id, policy)))
119            })
120            .collect::<Result<Vec<_>, _>>()?
121            .into_iter()
122            .flatten()
123            .collect();
124        let actions = registry
125            .capability_manifests()
126            .filter(|manifest| manifest.kind == CapabilityKind::Action)
127            .map(|manifest| (manifest.id.clone(), manifest.clone()))
128            .collect();
129        Ok(Self {
130            spec_id: spec.spec_id.clone(),
131            host,
132            schedules,
133            actions,
134        })
135    }
136
137    fn action_intents(
138        &self,
139        item: &WorkItem,
140        requests: Vec<crate::PreparedAction>,
141    ) -> Result<Vec<ActionIntent>, String> {
142        requests
143            .into_iter()
144            .enumerate()
145            .map(|(index, request)| {
146                request.validate().map_err(|error| error.to_string())?;
147                let manifest = self
148                    .actions
149                    .get(request.capability_id.as_str())
150                    .ok_or_else(|| {
151                        format!(
152                            "action step emitted unknown capability '{}'",
153                            request.capability_id
154                        )
155                    })?;
156                let capability = item
157                    .capability_pins
158                    .iter()
159                    .find(|pin| {
160                        pin.id == manifest.id
161                            && pin.contract_version == manifest.contract_version
162                            && pin.content_digest == manifest.content_digest
163                    })
164                    .cloned()
165                    .ok_or_else(|| {
166                        format!(
167                            "workflow revision does not pin action capability '{}'",
168                            manifest.id
169                        )
170                    })?;
171                Ok(ActionIntent {
172                    id: format!("{}:{}:action:{index}", item.id, item.state_version),
173                    tenant_id: item.tenant_id.clone(),
174                    instance_id: item
175                        .id
176                        .parse()
177                        .map_err(|error: af_context::EmptyId| error.to_string())?,
178                    run_id: item.run_id.clone(),
179                    capability,
180                    idempotency_key: format!(
181                        "{}:{}:{}",
182                        item.id, item.state_version, request.idempotency_key
183                    ),
184                    state: ActionState::Prepared,
185                    input: request.input,
186                    effect: manifest.effect,
187                    retry_class: manifest.idempotency_mode,
188                    control_epochs: item.control_epochs,
189                    resource_scope_id: request.resource_scope_id,
190                    lease_epoch: item.lease_version,
191                    action_epoch: item.state_version,
192                    // Provider timeout is per attempt. A product may set a
193                    // separate whole-intent deadline explicitly.
194                    deadline: request.deadline,
195                    reservation: request.reservation,
196                    created_at: item.claimed_at,
197                })
198            })
199            .collect()
200    }
201
202    fn terminal_action(
203        &self,
204        item: &WorkItem,
205        next_state: &mut Map<String, Value>,
206    ) -> Result<Option<WorkflowTransitionCommand>, String> {
207        let Some(wakeup) = item
208            .wakeups
209            .iter()
210            .find(|wakeup| wakeup.kind == "timer" && wakeup.payload["kind"] == "terminal_action")
211        else {
212            return Ok(None);
213        };
214        let observation: crate::ActionObservation =
215            serde_json::from_value(wakeup.payload["observation"].clone())
216                .map_err(|error| format!("terminal action fact: {error}"))?;
217        let Some(internal) = next_state
218            .get_mut("__workflow")
219            .and_then(Value::as_object_mut)
220        else {
221            return Err("terminal action has no durable pending-action state".into());
222        };
223        let pending_empty = {
224            let pending = internal
225                .get_mut("pending_actions")
226                .and_then(Value::as_array_mut)
227                .ok_or_else(|| "terminal action has no pending action list".to_string())?;
228            let before = pending.len();
229            pending.retain(|id| id.as_str() != Some(&observation.action_intent_id));
230            if pending.len() == before {
231                return Err("terminal action does not belong to this workflow instance".into());
232            }
233            pending.is_empty()
234        };
235        let succeeded = matches!(
236            observation.state.as_str(),
237            "completed" | "max_steps_reached" | "succeeded"
238        );
239        if !succeeded {
240            internal.insert("action_failed".into(), json!(true));
241        }
242        let failed = internal
243            .get("action_failed")
244            .and_then(Value::as_bool)
245            .unwrap_or(false);
246        if succeeded {
247            let result = observation.resource_ref.clone().unwrap_or(Value::Null);
248            if let Some(continuation) = internal
249                .get_mut("continuation")
250                .and_then(Value::as_object_mut)
251            {
252                let index = continuation
253                    .get("action_ids")
254                    .and_then(Value::as_array)
255                    .and_then(|ids| {
256                        ids.iter()
257                            .position(|id| id.as_str() == Some(&observation.action_intent_id))
258                    });
259                if let Some(index) = index {
260                    if let Some(event) = continuation
261                        .get_mut("events")
262                        .and_then(Value::as_array_mut)
263                        .and_then(|events| events.get_mut(index))
264                    {
265                        event["payload"]["action_result"] = result;
266                    }
267                } else if pending_empty {
268                    // Backward compatibility for single-action continuations
269                    // persisted before action_ids were recorded.
270                    if let Some(event) = continuation
271                        .get_mut("events")
272                        .and_then(Value::as_array_mut)
273                        .and_then(|events| (events.len() == 1).then(|| &mut events[0]))
274                    {
275                        event["payload"]["action_result"] = result;
276                    }
277                }
278            }
279        }
280        if pending_empty {
281            internal.remove("action_failed");
282            if failed {
283                internal.remove("continuation");
284            }
285        }
286        let resumes = pending_empty && !failed && internal.contains_key("continuation");
287        let disposition = if resumes {
288            WorkDisposition::Continue { delay_secs: 0 }
289        } else if pending_empty {
290            internal
291                .remove("resume_at")
292                .map(serde_json::from_value)
293                .transpose()
294                .map_err(|error| format!("stored workflow resume_at: {error}"))?
295                .map(|at| WorkDisposition::Reschedule { at })
296                .unwrap_or_else(|| {
297                    if internal
298                        .get("static_event_consumed")
299                        .and_then(Value::as_bool)
300                        .unwrap_or(false)
301                    {
302                        WorkDisposition::Complete
303                    } else {
304                        Self::idle_disposition()
305                    }
306                })
307        } else {
308            WorkDisposition::Continue {
309                delay_secs: IDLE_WAIT_SECS,
310            }
311        };
312        let reports_workflow_outcome = pending_empty && !resumes;
313        let mut selected = item.clone();
314        selected.wakeups = vec![wakeup.clone()];
315        Ok(Some(command(
316            &selected,
317            Value::Object(next_state.clone()),
318            Vec::new(),
319            EvaluationOutcome {
320                triggered: true,
321                matched: reports_workflow_outcome,
322                succeeded: succeeded && reports_workflow_outcome,
323                action_terminal: reports_workflow_outcome,
324            },
325            disposition,
326            "workflow.action_observed",
327        )))
328    }
329
330    fn catch_up_count(state: &Map<String, Value>) -> u32 {
331        state
332            .get("__workflow")
333            .and_then(Value::as_object)
334            .and_then(|internal| internal.get("catch_up_count"))
335            .and_then(Value::as_u64)
336            .and_then(|value| u32::try_from(value).ok())
337            .unwrap_or(0)
338    }
339
340    fn set_internal(state: &mut Map<String, Value>, key: &str, value: Value) -> Result<(), String> {
341        let internal = state
342            .entry("__workflow")
343            .or_insert_with(|| json!({}))
344            .as_object_mut()
345            .ok_or_else(|| "workflow internal state must be an object".to_string())?;
346        internal.insert(key.into(), value);
347        Ok(())
348    }
349
350    fn schedule_cursors(
351        &self,
352        item: &WorkItem,
353        state: &Map<String, Value>,
354    ) -> Result<BTreeMap<String, ScheduleCursor>, String> {
355        let saved = state
356            .get("__workflow")
357            .and_then(|value| value.get("schedules"));
358        let mut cursors: BTreeMap<String, ScheduleCursor> = saved
359            .cloned()
360            .map(serde_json::from_value)
361            .transpose()
362            .map_err(|error| format!("stored branch schedule: {error}"))?
363            .unwrap_or_default();
364        for (branch, policy) in &self.schedules {
365            if !cursors.contains_key(branch) {
366                if saved.is_some() {
367                    return Err("stored schedules omit a pinned branch".into());
368                }
369                // Preserve the previously executed root, including a pending action's resume cursor.
370                let legacy_root = self.host.branch_count() == 1
371                    || (saved.is_none()
372                        && item.state_version > 0
373                        && branch == self.host.branch_id()
374                        && (state.contains_key("last_run")
375                            || state.get("__workflow").is_some_and(|internal| {
376                                internal.get("catch_up_count").is_some()
377                                    || internal.get("resume_at").is_some()
378                            })));
379                let next_at = if legacy_root {
380                    state
381                        .get("__workflow")
382                        .and_then(|internal| internal.get("resume_at"))
383                        .cloned()
384                        .map(serde_json::from_value)
385                        .transpose()
386                        .map_err(|error| format!("stored workflow resume_at: {error}"))?
387                        .unwrap_or(item.scheduled_at)
388                } else {
389                    crate::next_scheduled_at(
390                        policy,
391                        item.lifecycle
392                            .starts_at
393                            .map_or(item.created_at, |starts| starts.max(item.created_at)),
394                        None,
395                    )?
396                };
397                cursors.insert(
398                    branch.clone(),
399                    ScheduleCursor {
400                        next_at,
401                        catch_up_count: if legacy_root {
402                            Self::catch_up_count(state)
403                        } else {
404                            0
405                        },
406                    },
407                );
408            }
409        }
410        if cursors
411            .keys()
412            .any(|branch| !self.schedules.contains_key(branch))
413        {
414            return Err("stored schedule refers to an unknown branch".into());
415        }
416        Ok(cursors)
417    }
418
419    fn idle_disposition() -> WorkDisposition {
420        WorkDisposition::Continue {
421            delay_secs: IDLE_WAIT_SECS,
422        }
423    }
424}
425
426#[async_trait]
427impl<Context: Send + Sync> WorkflowDriver<Context> for SpecDriver {
428    fn name(&self) -> &'static str {
429        "spec"
430    }
431
432    fn spec_ids(&self) -> Vec<&str> {
433        vec![&self.spec_id]
434    }
435
436    fn validate_specs(&self) -> Result<(), String> {
437        Ok(())
438    }
439
440    async fn evaluate(
441        &self,
442        _: &Context,
443        item: &WorkItem,
444    ) -> Result<WorkflowTransitionCommand, String> {
445        if item.spec_id != self.spec_id {
446            return Err("compiled spec does not match the claimed spec identity".into());
447        }
448        let mut next_state = match &item.config {
449            Value::Object(map) => map.clone(),
450            Value::Null => Map::new(),
451            _ => return Err("spec instance config must be a JSON object".into()),
452        };
453        if item.cancel_requested {
454            return Ok(command(
455                item,
456                Value::Object(next_state),
457                Vec::new(),
458                EvaluationOutcome::default(),
459                WorkDisposition::Complete,
460                "workflow.cancelled",
461            ));
462        }
463        if let Some(command) = self.terminal_action(item, &mut next_state)? {
464            return Ok(command);
465        }
466        if next_state
467            .get("__workflow")
468            .and_then(Value::as_object)
469            .and_then(|internal| internal.get("pending_actions"))
470            .and_then(Value::as_array)
471            .is_some_and(|pending| !pending.is_empty())
472        {
473            return Ok(command(
474                item,
475                Value::Object(next_state),
476                Vec::new(),
477                EvaluationOutcome::default(),
478                Self::idle_disposition(),
479                "workflow.action_waiting",
480            ));
481        }
482        let continuation = next_state
483            .get_mut("__workflow")
484            .and_then(Value::as_object_mut)
485            .and_then(|internal| internal.remove("continuation"))
486            .map(serde_json::from_value::<ActionContinuation>)
487            .transpose()
488            .map_err(|error| format!("stored workflow continuation: {error}"))?;
489        let resumed_from_actions = continuation.is_some();
490        let mut cursors = self.schedule_cursors(item, &next_state)?;
491        // ponytail: one transition per instance; alternate due schedules and deliveries
492        // under backlog. Use weighted fairness only if consumers need unequal shares.
493        let due = cursors
494            .values()
495            .any(|cursor| cursor.next_at <= item.claimed_at);
496        let prefer_schedule = due
497            && !item.wakeups.is_empty()
498            && next_state
499                .get("__workflow")
500                .and_then(|internal| internal.get("last_dispatch"))
501                .and_then(Value::as_str)
502                != Some("schedule");
503        let mut selected = item.clone();
504        if prefer_schedule {
505            selected.wakeups.clear();
506        }
507        let item = &selected;
508        let earliest = cursors
509            .iter()
510            .min_by_key(|(branch, cursor)| (cursor.next_at, *branch));
511        let branch_id = if let Some(continuation) = &continuation {
512            continuation.branch_id.clone()
513        } else if let Some(wakeup) = item.wakeups.first() {
514            wakeup
515                .branch_id
516                .as_ref()
517                .map(|id| id.as_str())
518                .or_else(|| (self.host.branch_count() == 1).then(|| self.host.branch_id()))
519                .ok_or_else(|| "multi-branch wakeup requires branch_id".to_string())?
520                .to_owned()
521        } else if let Some((branch, _)) = earliest {
522            branch.clone()
523        } else if self.host.branch_count() == 1 {
524            self.host.branch_id().to_owned()
525        } else {
526            return Err("multi-branch workflow requires a branch-addressed delivery".into());
527        };
528        if !self.host.branch_ids().any(|branch| branch == branch_id) {
529            return Err("wakeup refers to an unknown branch".into());
530        }
531        let schedule = if continuation.is_none() && item.wakeups.is_empty() {
532            if let Some(cursor) = cursors.get_mut(&branch_id) {
533                let decision = if cursor.next_at > item.claimed_at {
534                    crate::ScheduleDecision {
535                        tick_at: None,
536                        next_at: cursor.next_at,
537                        catch_up_count: cursor.catch_up_count,
538                    }
539                } else {
540                    schedule_decision(
541                        &self.schedules[&branch_id],
542                        cursor.next_at,
543                        item.claimed_at,
544                        cursor.catch_up_count,
545                    )?
546                };
547                cursor.next_at = decision.next_at;
548                cursor.catch_up_count = decision.catch_up_count;
549                Some(decision)
550            } else {
551                None
552            }
553        } else {
554            None
555        };
556        Self::set_internal(
557            &mut next_state,
558            "last_dispatch",
559            json!(if continuation.is_some() {
560                "continuation"
561            } else if item.wakeups.is_empty() {
562                "schedule"
563            } else {
564                "delivery"
565            }),
566        )?;
567        Self::set_internal(&mut next_state, "schedules", json!(cursors))?;
568        let next_at = cursors.values().map(|cursor| cursor.next_at).min();
569        if let Some(decision) = schedule {
570            if decision.tick_at.is_none() {
571                return Ok(command(
572                    item,
573                    Value::Object(next_state),
574                    Vec::new(),
575                    EvaluationOutcome::default(),
576                    WorkDisposition::Reschedule {
577                        at: next_at.expect("selected schedule has a cursor"),
578                    },
579                    "workflow.schedule_skipped",
580                ));
581            }
582        }
583        let state = Arc::new(MemoryState::from_snapshot(
584            next_state
585                .get("state")
586                .and_then(Value::as_object)
587                .cloned()
588                .unwrap_or_default(),
589        ));
590        let wakeup_payload = item.wakeups.first().map(|wakeup| wakeup.payload.clone());
591        let static_event = next_state.get("event").cloned().filter(|_| {
592            !next_state
593                .get("__workflow")
594                .and_then(Value::as_object)
595                .and_then(|internal| internal.get("static_event_consumed"))
596                .and_then(Value::as_bool)
597                .unwrap_or(false)
598        });
599        let consumed_static_event = wakeup_payload.is_none() && static_event.is_some();
600        let payload = wakeup_payload
601            .or(static_event)
602            .or_else(|| {
603                schedule
604                    .and_then(|decision| decision.tick_at)
605                    .map(|tick_at| json!({ "tick_at": tick_at }))
606            })
607            .or_else(|| continuation.as_ref().map(|_| json!({})))
608            .ok_or_else(|| "event workflow requires a trigger delivery".to_string())?;
609        let mut action_input = item.config.as_object().cloned().unwrap_or_default();
610        for internal in ["__workflow", "state", "last_run", "event"] {
611            action_input.remove(internal);
612        }
613        if let Some(trigger) = payload.as_object() {
614            action_input.extend(trigger.clone());
615        }
616        if consumed_static_event {
617            Self::set_internal(&mut next_state, "static_event_consumed", json!(true))?;
618        }
619        let context = self.host.context_for(&branch_id, state.clone());
620        let (outcome, continuation_step) = if let Some(continuation) = continuation {
621            let mut events = continuation.events.into_iter();
622            let first = events
623                .next()
624                .ok_or_else(|| "stored workflow continuation is empty".to_string())?;
625            let (mut combined, mut next_step) = self
626                .host
627                .run_event_with_continuation_on(
628                    &continuation.branch_id,
629                    &context,
630                    first,
631                    continuation.step,
632                )
633                .await
634                .ok_or_else(|| "stored workflow continuation branch is missing".to_string())?;
635            for event in events {
636                let (outcome, event_next_step) = self
637                    .host
638                    .run_event_with_continuation_on(
639                        &continuation.branch_id,
640                        &context,
641                        event,
642                        continuation.step,
643                    )
644                    .await
645                    .ok_or_else(|| "stored workflow continuation branch is missing".to_string())?;
646                if next_step.is_some() && event_next_step.is_some() && next_step != event_next_step
647                {
648                    return Err("continued events stopped at different action steps".into());
649                }
650                next_step = next_step.or(event_next_step);
651                combined.steps_run += outcome.steps_run;
652                combined.survivors.extend(outcome.survivors);
653                combined.actions.extend(outcome.actions);
654                combined.matched |= outcome.matched;
655                combined.succeeded &= outcome.succeeded;
656                if !combined.survivors.is_empty() || !combined.actions.is_empty() {
657                    combined.terminal = Terminal::Completed;
658                }
659            }
660            (combined, next_step)
661        } else {
662            self.host
663                .run_event_with_continuation_on(
664                    &branch_id,
665                    &context,
666                    Event::from_json(Value::Object(action_input)),
667                    0,
668                )
669                .await
670                .ok_or_else(|| "spec has no selected branch".to_string())?
671        };
672        let (terminal, exit_reason) = match &outcome.terminal {
673            Terminal::Completed => ("completed", Value::Null),
674            Terminal::Dropped { node_id, reason } => {
675                ("dropped", json!({ "node_id": node_id, "reason": reason }))
676            }
677        };
678        next_state.insert("state".into(), Value::Object(state.snapshot()));
679        next_state.insert(
680            "last_run".into(),
681            json!({
682                "branch_id": branch_id,
683                "terminal": terminal,
684                "exit": exit_reason,
685                "steps_run": outcome.steps_run,
686                "survivors": outcome.survivors.len(),
687                "at": item.claimed_at,
688            }),
689        );
690        let action_intents = self.action_intents(item, outcome.actions)?;
691        if let Some(step) = continuation_step {
692            if outcome.survivors.is_empty() || outcome.survivors.len() != action_intents.len() {
693                return Err("action continuation requires one surviving event per action".into());
694            }
695            Self::set_internal(
696                &mut next_state,
697                "continuation",
698                serde_json::to_value(ActionContinuation {
699                    branch_id: branch_id.clone(),
700                    step,
701                    events: outcome.survivors.clone(),
702                    action_ids: action_intents
703                        .iter()
704                        .map(|intent| intent.id.clone())
705                        .collect(),
706                })
707                .map_err(|error| error.to_string())?,
708            )?;
709        }
710        let disposition = match next_at {
711            Some(at) if action_intents.is_empty() => WorkDisposition::Reschedule { at },
712            Some(at) => {
713                Self::set_internal(&mut next_state, "resume_at", json!(at))?;
714                Self::idle_disposition()
715            }
716            None if action_intents.is_empty() && consumed_static_event => WorkDisposition::Complete,
717            None => Self::idle_disposition(),
718        };
719        if !action_intents.is_empty() {
720            let internal = next_state
721                .entry("__workflow")
722                .or_insert_with(|| json!({}))
723                .as_object_mut()
724                .ok_or_else(|| "workflow internal state must be an object".to_string())?;
725            let pending = internal
726                .entry("pending_actions")
727                .or_insert_with(|| json!([]))
728                .as_array_mut()
729                .ok_or_else(|| "workflow pending actions must be an array".to_string())?;
730            pending.extend(action_intents.iter().map(|intent| json!(intent.id)));
731        }
732        let evaluation = EvaluationOutcome {
733            triggered: true,
734            matched: outcome.matched,
735            succeeded: outcome.succeeded && action_intents.is_empty(),
736            action_terminal: resumed_from_actions && action_intents.is_empty(),
737        };
738        Ok(command(
739            item,
740            Value::Object(next_state),
741            action_intents,
742            evaluation,
743            disposition,
744            "workflow.spec_evaluated",
745        ))
746    }
747}
748
749fn command(
750    item: &WorkItem,
751    next_state: Value,
752    action_intents: Vec<ActionIntent>,
753    outcome: EvaluationOutcome,
754    disposition: WorkDisposition,
755    event_type: &str,
756) -> WorkflowTransitionCommand {
757    WorkflowTransitionCommand {
758        consumed_wakeups: if outcome.triggered {
759            item.wakeups
760                .first()
761                .map(|wakeup| wakeup.id.clone())
762                .into_iter()
763                .collect()
764        } else {
765            Vec::new()
766        },
767        delivery_key: format!("spec:{}:{}", item.id, item.state_version),
768        delivery_digest: format!(
769            "{}:{}:{}",
770            item.workflow_revision_digest, item.execution_profile_digest, item.state_version
771        ),
772        event_type: event_type.into(),
773        event_digest: format!("{event_type}:{}:{}", item.id, item.state_version),
774        event_payload: json!({ "spec_id": item.spec_id, "wakeups": item.wakeups.len() }),
775        next_state,
776        action_intents,
777        outcome,
778        disposition,
779    }
780}
781
782#[cfg(test)]
783mod tests {
784    use super::*;
785    use crate::{ControlEpochs, DriverRegistry, PreparedAction, StepNode, StepResult, Wakeup};
786
787    struct PassNode;
788
789    #[async_trait::async_trait]
790    impl StepNode for PassNode {
791        async fn process(&self, event: &Event, _: &crate::WorkflowContext) -> StepResult {
792            StepResult::Pass(event.clone())
793        }
794    }
795
796    struct ActionNode;
797
798    struct FanOutTwo;
799
800    #[async_trait::async_trait]
801    impl StepNode for FanOutTwo {
802        fn produces_fan_out(&self) -> bool {
803            true
804        }
805
806        async fn process(&self, event: &Event, _: &crate::WorkflowContext) -> StepResult {
807            StepResult::FanOut(vec![
808                Event::from_json(json!({"value":event.payload["value"]})),
809                Event::from_json(json!({"value":2})),
810            ])
811        }
812    }
813
814    #[async_trait::async_trait]
815    impl StepNode for ActionNode {
816        async fn process(&self, event: &Event, _: &crate::WorkflowContext) -> StepResult {
817            StepResult::Action {
818                event: event.clone(),
819                action: Box::new(PreparedAction::new(
820                    "execute.demo".parse().unwrap(),
821                    "request-1",
822                    Value::Object(event.payload.clone()),
823                )),
824            }
825        }
826    }
827
828    struct VersionNode(&'static str);
829
830    #[async_trait::async_trait]
831    impl StepNode for VersionNode {
832        async fn process(&self, event: &Event, ctx: &crate::WorkflowContext) -> StepResult {
833            ctx.state
834                .append(&ctx.scoped("seen"), json!(self.0), Some(3), None)
835                .await;
836            StepResult::Pass(event.clone())
837        }
838    }
839
840    fn version_one(_: &Value) -> Result<Box<dyn StepNode>, crate::NodeError> {
841        Ok(Box::new(VersionNode("v1")))
842    }
843
844    fn version_two(_: &Value) -> Result<Box<dyn StepNode>, crate::NodeError> {
845        Ok(Box::new(VersionNode("v2")))
846    }
847
848    fn item(spec_id: &str, config: Value) -> WorkItem {
849        WorkItem {
850            id: "instance".into(),
851            run_id: uuid::Uuid::new_v4().to_string().parse().unwrap(),
852            tenant_id: "tenant".parse().unwrap(),
853            subject_id: "subject".parse().unwrap(),
854            spec_id: spec_id.into(),
855            definition_id: spec_id.into(),
856            workflow_revision: 1,
857            workflow_revision_digest: "digest".into(),
858            execution_profile_id: "profile".into(),
859            execution_profile_revision: 1,
860            execution_profile_digest: "profile-digest".into(),
861            kernel_abi_version: "1".into(),
862            capability_pins: Vec::new(),
863            lifecycle: crate::LifecyclePolicy::run_once(),
864            scheduled_at: chrono::Utc::now(),
865            created_at: chrono::Utc::now(),
866            claimed_at: chrono::Utc::now(),
867            config,
868            state_version: 0,
869            control_epochs: ControlEpochs::default(),
870            cancel_requested: false,
871            lease_version: 1,
872            wakeups: Vec::new(),
873        }
874    }
875
876    fn spec(ingress: &str, ingress_config: Value) -> Spec {
877        Spec::from_json(
878            &json!({
879                "spec_id": "counter", "version": "1",
880                "branches": [{
881                    "branch_id": "__root__",
882                    "nodes": [
883                        {"id": "in", "type": ingress, "config": ingress_config},
884                        {"id": "count", "type": "transform.state_append",
885                         "config": {"key": "seen", "path": "value", "max_len": 3}}
886                    ],
887                    "edges": [{"source": "in", "target": "count"}]
888                }]
889            })
890            .to_string(),
891        )
892        .unwrap()
893    }
894
895    #[tokio::test]
896    async fn published_revisions_select_distinct_graphs_and_reject_conflicting_pins() {
897        let nodes = NodeRegistry::with_builtins();
898        let mut registry = crate::DriverRegistry::<()>::new().with_node_registry(nodes);
899        let first_spec = spec("ingress.event", json!({}));
900        registry.prewarm_spec(&first_spec).unwrap();
901        let first = item("counter", json!({"event": {"value": 7}}));
902        let revision = crate::WorkflowRevision {
903            definition_id: first.definition_id.clone(),
904            revision: 1,
905            content_digest: first.workflow_revision_digest.clone(),
906            kernel_abi_version: "1".into(),
907            dependency_set_digest: "deps".into(),
908            expression_versions: Default::default(),
909            capabilities: Vec::new(),
910            template_provenance: Value::Null,
911            spec: first_spec,
912        };
913        let first_driver = registry
914            .resolve_revision(&first, revision.clone())
915            .await
916            .unwrap();
917        let cached = registry
918            .resolve_revision(&first, revision.clone())
919            .await
920            .unwrap();
921        assert!(Arc::ptr_eq(&first_driver, &cached));
922        let mut second = first.clone();
923        second.workflow_revision = 2;
924        second.workflow_revision_digest = "revision-2".into();
925        let mut second_revision = revision.clone();
926        second_revision.revision = 2;
927        second_revision.content_digest = second.workflow_revision_digest.clone();
928        second_revision.spec.branches[0].nodes[1].config["key"] = json!("second");
929        let second_driver = registry
930            .resolve_revision(&second, second_revision.clone())
931            .await
932            .unwrap();
933        let first_result = first_driver.evaluate(&(), &first).await.unwrap();
934        let second_result = second_driver.evaluate(&(), &second).await.unwrap();
935        assert_eq!(
936            first_result.next_state["state"]["__root__.seen"],
937            json!([7])
938        );
939        assert_eq!(
940            second_result.next_state["state"]["__root__.second"],
941            json!([7])
942        );
943        assert!(second_result.next_state["state"]
944            .get("__root__.seen")
945            .is_none());
946        assert!(registry
947            .resolve_revision(&first, second_revision.clone())
948            .await
949            .is_err());
950        second_revision.spec.branches[0].nodes[1].config["key"] = json!("conflict");
951        assert!(registry
952            .resolve_revision(&second, second_revision.clone())
953            .await
954            .is_err());
955        let mut other_tenant = second.clone();
956        other_tenant.tenant_id = "other-tenant".parse().unwrap();
957        assert!(registry
958            .resolve_revision(&other_tenant, second_revision.clone())
959            .await
960            .is_ok());
961        let mut other_definition = second.clone();
962        other_definition.definition_id = "other-definition".into();
963        second_revision.definition_id = other_definition.definition_id.clone();
964        assert!(registry
965            .resolve_revision(&other_definition, second_revision)
966            .await
967            .is_ok());
968        let mut unsupported = revision.clone();
969        unsupported
970            .expression_versions
971            .insert("transform.state_append".into(), "unknown".into());
972        assert!(registry
973            .resolve_revision(&first, unsupported)
974            .await
975            .is_err());
976        let mut guarded_nodes = NodeRegistry::with_builtins();
977        let mut manifest = CapabilityManifest::action(
978            "transform.state_append",
979            "1",
980            "read-v1",
981            crate::Effect::Read,
982            crate::IdempotencyMode::None,
983            false,
984        );
985        manifest.kind = CapabilityKind::Expression;
986        guarded_nodes.register_capability(manifest).unwrap();
987        let guarded = crate::DriverRegistry::<()>::new().with_node_registry(guarded_nodes);
988        assert!(
989            guarded
990                .resolve_revision(&first, revision.clone())
991                .await
992                .is_err(),
993            "read capabilities need pins even though they do not emit actions"
994        );
995        for version in 3..=259 {
996            let mut next_item = first.clone();
997            next_item.workflow_revision = version;
998            next_item.workflow_revision_digest = format!("revision-{version}");
999            let mut next_revision = revision.clone();
1000            next_revision.revision = version;
1001            next_revision.content_digest = next_item.workflow_revision_digest.clone();
1002            registry
1003                .resolve_revision(&next_item, next_revision)
1004                .await
1005                .unwrap();
1006        }
1007        let (left, right) = tokio::join!(
1008            registry.resolve_revision(&first, revision.clone()),
1009            registry.resolve_revision(&first, revision.clone())
1010        );
1011        let (left, right) = (left.unwrap(), right.unwrap());
1012        assert!(
1013            !Arc::ptr_eq(&left, &first_driver),
1014            "oldest compiled revision is evicted"
1015        );
1016        assert!(
1017            Arc::ptr_eq(&left, &right),
1018            "concurrent resolution shares one compilation"
1019        );
1020        let replaced = registry.with_node_registry(NodeRegistry::empty());
1021        assert!(
1022            replaced
1023                .resolve_revision(&first, revision.clone())
1024                .await
1025                .is_err(),
1026            "changing the node registry must invalidate compiled graphs"
1027        );
1028        let mut unavailable = revision;
1029        unavailable.capabilities.push(crate::CapabilityPin {
1030            id: "missing".into(),
1031            contract_version: "1".into(),
1032            content_digest: "missing-v1".into(),
1033        });
1034        let mut unavailable_item = first;
1035        unavailable_item.capability_pins = unavailable.capabilities.clone();
1036        assert!(replaced
1037            .resolve_revision(&unavailable_item, unavailable)
1038            .await
1039            .is_err());
1040    }
1041
1042    #[tokio::test]
1043    async fn published_revisions_keep_two_versions_of_one_capability_runnable() {
1044        let capability = |version: &str, digest: &str| {
1045            let mut manifest = CapabilityManifest::action(
1046                "transform.state_append",
1047                version,
1048                digest,
1049                crate::Effect::Read,
1050                crate::IdempotencyMode::None,
1051                false,
1052            );
1053            manifest.kind = CapabilityKind::Expression;
1054            manifest
1055        };
1056        let v1 = capability("1", "state-append-v1");
1057        let v2 = capability("2", "state-append-v2");
1058        let pin = |manifest: &CapabilityManifest| crate::CapabilityPin {
1059            id: manifest.id.clone(),
1060            contract_version: manifest.contract_version.clone(),
1061            content_digest: manifest.content_digest.clone(),
1062        };
1063        let mut nodes = NodeRegistry::with_builtins();
1064        let schema = nodes.schema("transform.state_append").cloned().unwrap();
1065        nodes.register_capability(v1.clone()).unwrap();
1066        nodes.register_capability(v2.clone()).unwrap();
1067        nodes
1068            .register_capability_implementation(
1069                pin(&v1),
1070                version_one,
1071                Some(schema.clone()),
1072                None,
1073                false,
1074            )
1075            .unwrap();
1076        nodes
1077            .register_capability_implementation(pin(&v2), version_two, Some(schema), None, false)
1078            .unwrap();
1079        let registry = crate::DriverRegistry::<()>::new().with_node_registry(nodes);
1080        let mut first = item("counter", json!({"event": {"value": 1}}));
1081        first.capability_pins = vec![pin(&v1)];
1082        let mut first_revision = crate::WorkflowRevision {
1083            definition_id: first.definition_id.clone(),
1084            revision: first.workflow_revision,
1085            content_digest: first.workflow_revision_digest.clone(),
1086            kernel_abi_version: "1".into(),
1087            dependency_set_digest: "deps-v1".into(),
1088            expression_versions: Default::default(),
1089            capabilities: first.capability_pins.clone(),
1090            template_provenance: Value::Null,
1091            spec: spec("ingress.event", json!({})),
1092        };
1093        let mut second = first.clone();
1094        second.workflow_revision = 2;
1095        second.workflow_revision_digest = "revision-2".into();
1096        second.capability_pins = vec![pin(&v2)];
1097        let mut second_revision = first_revision.clone();
1098        second_revision.revision = 2;
1099        second_revision.content_digest = second.workflow_revision_digest.clone();
1100        second_revision.dependency_set_digest = "deps-v2".into();
1101        second_revision.capabilities = second.capability_pins.clone();
1102
1103        let first_driver = registry
1104            .resolve_revision(&first, first_revision.clone())
1105            .await
1106            .unwrap();
1107        let second_driver = registry
1108            .resolve_revision(&second, second_revision)
1109            .await
1110            .unwrap();
1111
1112        assert_eq!(
1113            first_driver.evaluate(&(), &first).await.unwrap().next_state["state"]["__root__.seen"],
1114            json!(["v1"])
1115        );
1116        assert_eq!(
1117            second_driver
1118                .evaluate(&(), &second)
1119                .await
1120                .unwrap()
1121                .next_state["state"]["__root__.seen"],
1122            json!(["v2"])
1123        );
1124        first_revision.capabilities = vec![pin(&v2)];
1125        assert!(registry
1126            .resolve_revision(&first, first_revision)
1127            .await
1128            .is_err());
1129    }
1130
1131    #[tokio::test]
1132    async fn old_root_cursor_is_preserved_and_due_schedules_are_not_starved() {
1133        let mut graph = spec("ingress.fixed_rate", json!({"milliseconds": 10_000}));
1134        let mut sibling = spec("ingress.event", json!({})).branches.remove(0);
1135        sibling.branch_id = "events".into();
1136        for node in &mut sibling.nodes {
1137            node.id = format!("events-{}", node.id);
1138        }
1139        for edge in &mut sibling.edges {
1140            edge.source = format!("events-{}", edge.source);
1141            edge.target = format!("events-{}", edge.target);
1142        }
1143        graph.branches.push(sibling);
1144        let driver = SpecDriver::new(&graph, &NodeRegistry::with_builtins()).unwrap();
1145        let start: chrono::DateTime<chrono::Utc> = "2026-01-01T00:00:00Z".parse().unwrap();
1146        let mut work = item("counter", json!({"__workflow":{"catch_up_count":2}}));
1147        work.created_at = start;
1148        work.state_version = 7;
1149        work.scheduled_at = start + chrono::Duration::seconds(100);
1150        work.claimed_at = work.scheduled_at;
1151        let cursors = driver
1152            .schedule_cursors(&work, work.config.as_object().unwrap())
1153            .unwrap();
1154        assert_eq!(cursors["__root__"].next_at, work.scheduled_at);
1155        assert_eq!(cursors["__root__"].catch_up_count, 2);
1156        work.config["__workflow"]["resume_at"] = json!(start + chrono::Duration::seconds(150));
1157        let cursors = driver
1158            .schedule_cursors(&work, work.config.as_object().unwrap())
1159            .unwrap();
1160        assert_eq!(
1161            cursors["__root__"].next_at,
1162            start + chrono::Duration::seconds(150)
1163        );
1164        work.config["__workflow"]
1165            .as_object_mut()
1166            .unwrap()
1167            .remove("resume_at");
1168        work.wakeups.push(serde_json::from_value(json!({"id":"queued", "kind":"delivery", "branch_id":"events", "payload":{"value":7}})).unwrap());
1169        let schedule = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
1170            .await
1171            .unwrap();
1172        assert_eq!(schedule.next_state["last_run"]["branch_id"], "__root__");
1173        assert!(
1174            schedule.consumed_wakeups.is_empty(),
1175            "schedule does not acknowledge queued event"
1176        );
1177        work.config = schedule.next_state;
1178        let event = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
1179            .await
1180            .unwrap();
1181        assert_eq!(event.next_state["last_run"]["branch_id"], "events");
1182        assert_eq!(event.consumed_wakeups, ["queued"]);
1183        work.config = event.next_state;
1184        work.claimed_at += chrono::Duration::seconds(10);
1185        let schedule = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
1186            .await
1187            .unwrap();
1188        assert_eq!(schedule.next_state["last_run"]["branch_id"], "__root__");
1189        assert!(schedule.consumed_wakeups.is_empty());
1190    }
1191
1192    #[test]
1193    fn first_failure_does_not_migrate_a_nonexistent_root_cursor() {
1194        let mut graph = spec("ingress.fixed_rate", json!({"milliseconds":100_000}));
1195        let mut fast = graph.branches[0].clone();
1196        fast.branch_id = "fast".into();
1197        fast.nodes[0].config["milliseconds"] = json!(10_000);
1198        for node in &mut fast.nodes {
1199            node.id = format!("fast-{}", node.id);
1200        }
1201        for edge in &mut fast.edges {
1202            edge.source = format!("fast-{}", edge.source);
1203            edge.target = format!("fast-{}", edge.target);
1204        }
1205        graph.branches.push(fast);
1206        let driver = SpecDriver::new(&graph, &NodeRegistry::with_builtins()).unwrap();
1207        let mut work = item("counter", json!({}));
1208        work.created_at = "2026-01-01T00:00:00Z".parse().unwrap();
1209        work.scheduled_at = work.created_at + chrono::Duration::seconds(10);
1210        work.state_version = 1;
1211        let cursors = driver
1212            .schedule_cursors(&work, work.config.as_object().unwrap())
1213            .unwrap();
1214        assert_eq!(
1215            cursors["__root__"].next_at,
1216            work.created_at + chrono::Duration::seconds(100)
1217        );
1218        assert_eq!(cursors["fast"].next_at, work.scheduled_at);
1219    }
1220
1221    #[tokio::test]
1222    async fn branches_keep_independent_schedule_cursors_across_restart() {
1223        let mut graph = spec("ingress.event", json!({}));
1224        for (branch_id, period) in [("fast", 10_000), ("slow", 25_000)] {
1225            let mut branch = spec("ingress.fixed_rate", json!({"milliseconds": period}))
1226                .branches
1227                .remove(0);
1228            branch.branch_id = branch_id.into();
1229            for node in &mut branch.nodes {
1230                node.id = format!("{branch_id}-{}", node.id);
1231            }
1232            for edge in &mut branch.edges {
1233                edge.source = format!("{branch_id}-{}", edge.source);
1234                edge.target = format!("{branch_id}-{}", edge.target);
1235            }
1236            branch.nodes[1].config["path"] = json!("tick_at");
1237            graph.branches.push(branch);
1238        }
1239        let start: chrono::DateTime<chrono::Utc> = "2026-01-01T00:00:00Z".parse().unwrap();
1240        let mut work = item("counter", json!({}));
1241        work.created_at = start;
1242        work.scheduled_at = start + chrono::Duration::seconds(10);
1243        for (second, branch) in [(10, "fast"), (20, "fast"), (25, "slow"), (30, "fast")] {
1244            work.claimed_at = start + chrono::Duration::seconds(second);
1245            let restarted = SpecDriver::new(&graph, &NodeRegistry::with_builtins()).unwrap();
1246            let command = WorkflowDriver::<()>::evaluate(&restarted, &(), &work)
1247                .await
1248                .unwrap();
1249            assert_eq!(command.next_state["last_run"]["branch_id"], branch);
1250            let WorkDisposition::Reschedule { at } = command.disposition else {
1251                panic!("scheduled branch must retain every cursor");
1252            };
1253            work.config = serde_json::from_str(&command.next_state.to_string()).unwrap();
1254            work.scheduled_at = at;
1255        }
1256        assert_eq!(
1257            work.config["state"]["fast.seen"].as_array().unwrap().len(),
1258            3
1259        );
1260        assert_eq!(
1261            work.config["state"]["slow.seen"],
1262            json!([start + chrono::Duration::seconds(25)])
1263        );
1264        let saved_cursors = work.config["__workflow"]["schedules"].clone();
1265        work.wakeups.push(serde_json::from_value(json!({"id":"event", "kind":"delivery", "branch_id":"__root__", "payload":{"value":99}})).unwrap());
1266        let driver = SpecDriver::new(&graph, &NodeRegistry::with_builtins()).unwrap();
1267        let command = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
1268            .await
1269            .unwrap();
1270        assert_eq!(command.next_state["__workflow"]["schedules"], saved_cursors);
1271        assert_eq!(command.next_state["state"]["__root__.seen"], json!([99]));
1272        work.wakeups.clear();
1273        work.claimed_at = start + chrono::Duration::seconds(31);
1274        let skipped = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
1275            .await
1276            .unwrap();
1277        assert_eq!(skipped.event_type, "workflow.schedule_skipped");
1278        assert_eq!(skipped.next_state["state"], work.config["state"]);
1279    }
1280
1281    #[tokio::test]
1282    async fn delivery_selects_branch_and_rejects_ambiguous_legacy_identity() {
1283        let mut graph = spec("ingress.event", json!({}));
1284        let mut sibling = graph.branches[0].clone();
1285        sibling.branch_id = "other".into();
1286        for node in &mut sibling.nodes {
1287            node.id = format!("other-{}", node.id);
1288        }
1289        for edge in &mut sibling.edges {
1290            edge.source = format!("other-{}", edge.source);
1291            edge.target = format!("other-{}", edge.target);
1292        }
1293        graph.branches.push(sibling);
1294        let driver = SpecDriver::new(&graph, &NodeRegistry::with_builtins()).unwrap();
1295        let mut work = item("counter", json!({}));
1296        work.wakeups.push(
1297            serde_json::from_value(json!({
1298                "id": "delivery", "kind": "delivery", "branch_id": "other",
1299                "payload": {"value": 7}
1300            }))
1301            .unwrap(),
1302        );
1303        let command = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
1304            .await
1305            .unwrap();
1306        assert_eq!(command.next_state["state"]["other.seen"], json!([7]));
1307        assert!(command.next_state["state"]["__root__.seen"].is_null());
1308        work.config = command.next_state;
1309        work.wakeups[0] = serde_json::from_value(json!({
1310            "id": "next", "kind": "delivery", "branch_id": "__root__",
1311            "payload": {"value": 9}
1312        }))
1313        .unwrap();
1314        let restarted = SpecDriver::new(&graph, &NodeRegistry::with_builtins()).unwrap();
1315        let command = WorkflowDriver::<()>::evaluate(&restarted, &(), &work)
1316            .await
1317            .unwrap();
1318        assert_eq!(command.next_state["state"]["other.seen"], json!([7]));
1319        assert_eq!(command.next_state["state"]["__root__.seen"], json!([9]));
1320        for branch in [Value::Null, json!("missing")] {
1321            work.wakeups[0] = serde_json::from_value(json!({
1322                "id": "bad", "kind": "delivery", "branch_id": branch,
1323                "payload": {"value": 1}
1324            }))
1325            .unwrap();
1326            assert!(WorkflowDriver::<()>::evaluate(&driver, &(), &work)
1327                .await
1328                .is_err());
1329        }
1330    }
1331
1332    #[tokio::test]
1333    async fn compiled_driver_rejects_another_spec_identity() {
1334        let registry = NodeRegistry::with_builtins();
1335        let driver = SpecDriver::new(&spec("ingress.event", json!({})), &registry).unwrap();
1336        let result = WorkflowDriver::<()>::evaluate(
1337            &driver,
1338            &(),
1339            &item("different-spec", json!({"event": {"value": 1}})),
1340        )
1341        .await;
1342        assert!(
1343            result.is_err(),
1344            "a compiled graph must not execute another spec"
1345        );
1346    }
1347
1348    #[tokio::test]
1349    async fn cron_spec_reschedules_and_persists_branch_state() {
1350        let registry = NodeRegistry::with_builtins();
1351        let driver = SpecDriver::new(
1352            &spec("ingress.cron", json!({"expression": "0 0 * * * *"})),
1353            &registry,
1354        )
1355        .unwrap();
1356        let first = WorkflowDriver::<()>::evaluate(
1357            &driver,
1358            &(),
1359            &item("counter", json!({"event": {"value": 1}})),
1360        )
1361        .await
1362        .unwrap();
1363        let WorkDisposition::Reschedule { at } = first.disposition else {
1364            panic!("cron spec must reschedule");
1365        };
1366        assert_eq!(first.next_state["last_run"]["terminal"], "completed");
1367        let mut next = item("counter", first.next_state);
1368        next.scheduled_at = at;
1369        next.claimed_at = at;
1370        let second = WorkflowDriver::<()>::evaluate(&driver, &(), &next)
1371            .await
1372            .unwrap();
1373        assert_eq!(
1374            second.next_state["state"]["__root__.seen"],
1375            json!([1]),
1376            "the consumed bootstrap event must not be replayed on a cron tick"
1377        );
1378        assert_eq!(second.next_state["last_run"]["at"], json!(at));
1379    }
1380
1381    #[tokio::test]
1382    async fn event_spec_consumes_a_wakeup_then_waits() {
1383        let registry = NodeRegistry::with_builtins();
1384        let driver = SpecDriver::new(&spec("ingress.event", json!({})), &registry).unwrap();
1385        assert_eq!(
1386            WorkflowDriver::<()>::evaluate(&driver, &(), &item("counter", json!({})))
1387                .await
1388                .unwrap_err(),
1389            "event workflow requires a trigger delivery"
1390        );
1391        let mut work = item("counter", json!({}));
1392        work.wakeups.push(Wakeup {
1393            id: "delivery".into(),
1394            branch_id: None,
1395            kind: "delivery".into(),
1396            payload: json!({"value": 7}),
1397        });
1398        let mut second_wakeup = work.wakeups[0].clone();
1399        second_wakeup.id = "later-delivery".into();
1400        second_wakeup.payload = json!({"value":8});
1401        work.wakeups.push(second_wakeup);
1402        let command = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
1403            .await
1404            .unwrap();
1405        assert_eq!(command.consumed_wakeups, ["delivery"]);
1406        assert_eq!(command.next_state["state"]["__root__.seen"], json!([7]));
1407        assert!(matches!(
1408            command.disposition,
1409            WorkDisposition::Continue { .. }
1410        ));
1411        let mut registry = DriverRegistry::<()>::new();
1412        registry.register(Arc::new(driver)).unwrap();
1413        assert_eq!(registry.spec_ids(), ["counter"]);
1414    }
1415
1416    #[tokio::test]
1417    async fn effectful_steps_are_lifted_into_pinned_intents() {
1418        let mut registry = NodeRegistry::with_builtins();
1419        registry.register_step("guard.auth", |_| Ok(Box::new(PassNode)));
1420        registry.register_side_effect_guard("guard.auth").unwrap();
1421        registry.register_step("execute.demo", |_| Ok(Box::new(ActionNode)));
1422        let manifest = CapabilityManifest::action(
1423            "execute.demo",
1424            "1",
1425            "demo-digest",
1426            crate::Effect::ExternalWrite,
1427            crate::IdempotencyMode::Native,
1428            true,
1429        );
1430        registry.register_capability(manifest.clone()).unwrap();
1431        let effect = Spec::from_json(
1432            &json!({
1433                "spec_id": "effect", "version": "1",
1434                "branches": [{
1435                    "branch_id": "__root__",
1436                        "nodes": [
1437                        {"id": "in", "type": "ingress.event", "config": {}},
1438                        {"id": "auth", "type": "guard.auth", "config": {}},
1439                        {"id": "do", "type": "execute.demo", "config": {}}
1440                    ],
1441                    "edges": [
1442                        {"source": "in", "target": "auth"},
1443                        {"source": "auth", "target": "do"}
1444                    ]
1445                }]
1446            })
1447            .to_string(),
1448        )
1449        .unwrap();
1450        let driver = SpecDriver::new(&effect, &registry).unwrap();
1451        let mut work = item(
1452            "effect",
1453            json!({"objective":"review","event": {"value": 1}}),
1454        );
1455        work.capability_pins.push(crate::CapabilityPin {
1456            id: manifest.id,
1457            contract_version: manifest.contract_version,
1458            content_digest: manifest.content_digest,
1459        });
1460        let command = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
1461            .await
1462            .unwrap();
1463        assert_eq!(command.action_intents.len(), 1);
1464        assert_eq!(command.action_intents[0].state, ActionState::Prepared);
1465        assert_eq!(command.action_intents[0].capability.id, "execute.demo");
1466        assert_eq!(
1467            command.action_intents[0].input,
1468            json!({"objective":"review","value":1})
1469        );
1470        assert_eq!(command.action_intents[0].deadline, None);
1471
1472        let mut waiting = item("effect", command.next_state.clone());
1473        waiting.capability_pins = work.capability_pins.clone();
1474        let waiting_command = WorkflowDriver::<()>::evaluate(&driver, &(), &waiting)
1475            .await
1476            .unwrap();
1477        assert!(waiting_command.action_intents.is_empty());
1478        assert_eq!(waiting_command.event_type, "workflow.action_waiting");
1479
1480        waiting.wakeups.push(Wakeup {
1481            id: "terminal".into(),
1482            branch_id: None,
1483            kind: "timer".into(),
1484            payload: json!({
1485                "kind": "terminal_action",
1486                "observation": {
1487                    "id": "observation",
1488                    "action_intent_id": command.action_intents[0].id,
1489                    "provider_version": "1",
1490                    "observed_at": chrono::Utc::now(),
1491                    "state": "rejected",
1492                    "resource_ref": null,
1493                    "raw_receipt_digest": "receipt",
1494                    "terminal": true,
1495                    "retry_authorized": false
1496                }
1497            }),
1498        });
1499        let mut forged = waiting.clone();
1500        forged.wakeups[0].kind = "delivery".into();
1501        assert!(!forged.is_converging());
1502        let rejected = WorkflowDriver::<()>::evaluate(&driver, &(), &forged)
1503            .await
1504            .unwrap();
1505        assert_eq!(rejected.event_type, "workflow.action_waiting");
1506        assert!(rejected.consumed_wakeups.is_empty());
1507        waiting.wakeups.insert(
1508            0,
1509            Wakeup {
1510                id: "unrelated".into(),
1511                kind: "delivery".into(),
1512                branch_id: None,
1513                payload: json!({"value":999}),
1514            },
1515        );
1516        let terminal = WorkflowDriver::<()>::evaluate(&driver, &(), &waiting)
1517            .await
1518            .unwrap();
1519        assert_eq!(terminal.consumed_wakeups, [waiting.wakeups[1].id.clone()]);
1520        assert!(terminal.action_intents.is_empty());
1521        assert!(matches!(terminal.disposition, WorkDisposition::Complete));
1522        assert!(matches!(
1523            SpecDriver::new(&spec("ingress.cron", json!({})), &registry),
1524            Err(HostError::Schedule { .. })
1525        ));
1526    }
1527
1528    #[tokio::test]
1529    async fn bounded_fanout_actions_resume_with_their_own_results() {
1530        let mut registry = NodeRegistry::with_builtins();
1531        registry.register_step("guard.auth", |_| Ok(Box::new(PassNode)));
1532        registry.register_side_effect_guard("guard.auth").unwrap();
1533        registry.register_step("map.batch", |_| Ok(Box::new(FanOutTwo)));
1534        registry.register_fan_out("map.batch");
1535        registry.register_step("execute.demo", |_| Ok(Box::new(ActionNode)));
1536        let manifest = CapabilityManifest::action(
1537            "execute.demo",
1538            "1",
1539            "demo-digest",
1540            crate::Effect::ExternalWrite,
1541            crate::IdempotencyMode::Native,
1542            true,
1543        );
1544        registry.register_capability(manifest.clone()).unwrap();
1545        let graph = Spec::from_json(&json!({
1546            "spec_id":"fanout-actions","version":"1","branches":[{
1547                "branch_id":"__root__","nodes":[
1548                    {"id":"in","type":"ingress.event","config":{}},
1549                    {"id":"auth","type":"guard.auth","config":{}},
1550                    {"id":"map","type":"map.batch","config":{"count":2}},
1551                    {"id":"action","type":"execute.demo","config":{}},
1552                    {"id":"done","type":"transform.state_append","config":{"key":"results","path":"action_result.result","max_len":2}}
1553                ],"edges":[
1554                    {"source":"in","target":"auth"},{"source":"auth","target":"map"},
1555                    {"source":"map","target":"action"},{"source":"action","target":"done"}
1556                ]
1557            }]
1558        }).to_string()).unwrap();
1559        let driver = SpecDriver::new(&graph, &registry).unwrap();
1560        let mut work = item("fanout-actions", json!({"event":{"value":1}}));
1561        work.capability_pins.push(crate::CapabilityPin {
1562            id: manifest.id,
1563            contract_version: manifest.contract_version,
1564            content_digest: manifest.content_digest,
1565        });
1566        let first = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
1567            .await
1568            .unwrap();
1569        assert_eq!(first.action_intents.len(), 2);
1570        let mut waiting = item("fanout-actions", first.next_state);
1571        waiting.capability_pins = work.capability_pins;
1572        for (index, intent) in first.action_intents.iter().enumerate() {
1573            waiting.wakeups = vec![Wakeup {
1574                id: format!("terminal-{index}"),
1575                branch_id: None,
1576                kind: "timer".into(),
1577                payload: json!({"kind":"terminal_action","observation":{
1578                    "id":format!("observation-{index}"),"action_intent_id":intent.id,
1579                    "provider_version":"1","observed_at":chrono::Utc::now(),
1580                    "state":"succeeded","resource_ref":{"result":format!("result-{index}")},
1581                    "raw_receipt_digest":format!("receipt-{index}"),"terminal":true,
1582                    "retry_authorized":false
1583                }}),
1584            }];
1585            let observed = WorkflowDriver::<()>::evaluate(&driver, &(), &waiting)
1586                .await
1587                .unwrap();
1588            assert!(!observed.outcome.matched);
1589            assert!(!observed.outcome.succeeded);
1590            assert!(!observed.outcome.action_terminal);
1591            waiting.config = observed.next_state;
1592        }
1593        waiting.wakeups.clear();
1594        let resumed = WorkflowDriver::<()>::evaluate(&driver, &(), &waiting)
1595            .await
1596            .unwrap();
1597        assert_eq!(
1598            resumed.next_state["state"]["__root__.results"],
1599            json!(["result-0", "result-1"])
1600        );
1601        assert!(resumed.outcome.matched);
1602        assert!(resumed.outcome.succeeded);
1603        assert!(resumed.outcome.action_terminal);
1604    }
1605
1606    #[tokio::test]
1607    async fn successful_actions_resume_in_order_and_failures_stop_the_chain() {
1608        let mut registry = NodeRegistry::with_builtins();
1609        registry.register_step("guard.auth", |_| Ok(Box::new(PassNode)));
1610        registry.register_side_effect_guard("guard.auth").unwrap();
1611        registry.register_step("execute.demo", |_| Ok(Box::new(ActionNode)));
1612        let manifest = CapabilityManifest::action(
1613            "execute.demo",
1614            "1",
1615            "demo-digest",
1616            crate::Effect::ExternalWrite,
1617            crate::IdempotencyMode::Native,
1618            true,
1619        );
1620        registry.register_capability(manifest.clone()).unwrap();
1621        let graph = Spec::from_json(
1622            &json!({
1623                "spec_id": "chain", "version": "1",
1624                "branches": [{
1625                    "branch_id": "__root__",
1626                    "nodes": [
1627                        {"id":"in","type":"ingress.event","config":{}},
1628                        {"id":"auth","type":"guard.auth","config":{}},
1629                        {"id":"a","type":"execute.demo","config":{}},
1630                        {"id":"b","type":"execute.demo","config":{}},
1631                        {"id":"done","type":"transform.state_append","config":{"key":"results","path":"action_result.result","max_len":3}}
1632                    ],
1633                    "edges": [
1634                        {"source":"in","target":"auth"},
1635                        {"source":"auth","target":"a"},
1636                        {"source":"a","target":"b"},
1637                        {"source":"b","target":"done"}
1638                    ]
1639                }]
1640            }).to_string(),
1641        ).unwrap();
1642        let driver = || SpecDriver::new(&graph, &registry).unwrap();
1643        let pin = crate::CapabilityPin {
1644            id: manifest.id.clone(),
1645            contract_version: manifest.contract_version.clone(),
1646            content_digest: manifest.content_digest.clone(),
1647        };
1648        let mut work = item("chain", json!({"event":{"value":1}}));
1649        work.capability_pins.push(pin.clone());
1650        let first = WorkflowDriver::<()>::evaluate(&driver(), &(), &work)
1651            .await
1652            .unwrap();
1653        assert_eq!(first.action_intents.len(), 1);
1654
1655        let terminal = |action_intent_id: &str, state: &str, result: Value| Wakeup {
1656            id: format!("terminal:{action_intent_id}"),
1657            branch_id: None,
1658            kind: "timer".into(),
1659            payload: json!({"kind":"terminal_action","observation":{
1660                "id":format!("observation:{action_intent_id}"),
1661                "action_intent_id":action_intent_id,
1662                "provider_version":"1","observed_at":chrono::Utc::now(),
1663                "state":state,"resource_ref":result,"raw_receipt_digest":"receipt",
1664                "terminal":true,"retry_authorized":false
1665            }}),
1666        };
1667        let mut observed_a = item("chain", first.next_state);
1668        observed_a.capability_pins.push(pin.clone());
1669        observed_a.state_version = 1;
1670        observed_a.wakeups.push(terminal(
1671            &first.action_intents[0].id,
1672            "succeeded",
1673            json!({"result":"A"}),
1674        ));
1675        let wake_b = WorkflowDriver::<()>::evaluate(&driver(), &(), &observed_a)
1676            .await
1677            .unwrap();
1678        let mut resume_b = item("chain", wake_b.next_state);
1679        resume_b.capability_pins.push(pin.clone());
1680        resume_b.state_version = 2;
1681        let second = WorkflowDriver::<()>::evaluate(&driver(), &(), &resume_b)
1682            .await
1683            .unwrap();
1684        assert_eq!(second.action_intents.len(), 1);
1685        assert_eq!(
1686            second.action_intents[0].input["action_result"]["result"],
1687            "A"
1688        );
1689
1690        let mut observed_b = item("chain", second.next_state.clone());
1691        observed_b.capability_pins.push(pin.clone());
1692        observed_b.state_version = 3;
1693        observed_b.wakeups.push(terminal(
1694            &second.action_intents[0].id,
1695            "succeeded",
1696            json!({"result":"B"}),
1697        ));
1698        let wake_done = WorkflowDriver::<()>::evaluate(&driver(), &(), &observed_b)
1699            .await
1700            .unwrap();
1701        let mut resume_done = item("chain", wake_done.next_state);
1702        resume_done.capability_pins.push(pin.clone());
1703        resume_done.state_version = 4;
1704        let done = WorkflowDriver::<()>::evaluate(&driver(), &(), &resume_done)
1705            .await
1706            .unwrap();
1707        assert_eq!(done.next_state["state"]["__root__.results"], json!(["B"]));
1708
1709        let mut failed = item("chain", second.next_state);
1710        failed.capability_pins.push(pin);
1711        failed.state_version = 3;
1712        failed.wakeups.push(terminal(
1713            &second.action_intents[0].id,
1714            "failed",
1715            Value::Null,
1716        ));
1717        let stopped = WorkflowDriver::<()>::evaluate(&driver(), &(), &failed)
1718            .await
1719            .unwrap();
1720        assert!(matches!(stopped.disposition, WorkDisposition::Complete));
1721        assert!(stopped.next_state["__workflow"]["continuation"].is_null());
1722    }
1723
1724    #[tokio::test]
1725    async fn delivery_does_not_advance_the_cron_cursor() {
1726        let registry = NodeRegistry::with_builtins();
1727        let driver = SpecDriver::new(
1728            &spec("ingress.cron", json!({"expression": "0 0 * * * *"})),
1729            &registry,
1730        )
1731        .unwrap();
1732        let mut work = item("counter", json!({}));
1733        work.scheduled_at = work.claimed_at + chrono::Duration::hours(1);
1734        work.wakeups.push(Wakeup {
1735            id: "delivery".into(),
1736            branch_id: None,
1737            kind: "delivery".into(),
1738            payload: json!({"value": 7}),
1739        });
1740        let command = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
1741            .await
1742            .unwrap();
1743        assert_eq!(command.next_state["state"]["__root__.seen"], json!([7]));
1744        assert!(matches!(
1745            command.disposition,
1746            WorkDisposition::Reschedule { at } if at == work.scheduled_at
1747        ));
1748    }
1749}