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
620            .host
621            .context_for(&branch_id, state.clone())
622            .with_config(item.config.clone());
623        let (outcome, continuation_step) = if let Some(continuation) = continuation {
624            let mut events = continuation.events.into_iter();
625            let first = events
626                .next()
627                .ok_or_else(|| "stored workflow continuation is empty".to_string())?;
628            let (mut combined, mut next_step) = self
629                .host
630                .run_event_with_continuation_on(
631                    &continuation.branch_id,
632                    &context,
633                    first,
634                    continuation.step,
635                )
636                .await
637                .ok_or_else(|| "stored workflow continuation branch is missing".to_string())?;
638            for event in events {
639                let (outcome, event_next_step) = self
640                    .host
641                    .run_event_with_continuation_on(
642                        &continuation.branch_id,
643                        &context,
644                        event,
645                        continuation.step,
646                    )
647                    .await
648                    .ok_or_else(|| "stored workflow continuation branch is missing".to_string())?;
649                if next_step.is_some() && event_next_step.is_some() && next_step != event_next_step
650                {
651                    return Err("continued events stopped at different action steps".into());
652                }
653                next_step = next_step.or(event_next_step);
654                combined.steps_run += outcome.steps_run;
655                combined.survivors.extend(outcome.survivors);
656                combined.actions.extend(outcome.actions);
657                combined.matched |= outcome.matched;
658                combined.succeeded &= outcome.succeeded;
659                if !combined.survivors.is_empty() || !combined.actions.is_empty() {
660                    combined.terminal = Terminal::Completed;
661                }
662            }
663            (combined, next_step)
664        } else {
665            self.host
666                .run_event_with_continuation_on(
667                    &branch_id,
668                    &context,
669                    Event::from_json(Value::Object(action_input)),
670                    0,
671                )
672                .await
673                .ok_or_else(|| "spec has no selected branch".to_string())?
674        };
675        let (terminal, exit_reason) = match &outcome.terminal {
676            Terminal::Completed => ("completed", Value::Null),
677            Terminal::Dropped { node_id, reason } => {
678                ("dropped", json!({ "node_id": node_id, "reason": reason }))
679            }
680        };
681        next_state.insert("state".into(), Value::Object(state.snapshot()));
682        next_state.insert(
683            "last_run".into(),
684            json!({
685                "branch_id": branch_id,
686                "terminal": terminal,
687                "exit": exit_reason,
688                "steps_run": outcome.steps_run,
689                "survivors": outcome.survivors.len(),
690                "at": item.claimed_at,
691            }),
692        );
693        let action_intents = self.action_intents(item, outcome.actions)?;
694        if let Some(step) = continuation_step {
695            if outcome.survivors.is_empty() || outcome.survivors.len() != action_intents.len() {
696                return Err("action continuation requires one surviving event per action".into());
697            }
698            Self::set_internal(
699                &mut next_state,
700                "continuation",
701                serde_json::to_value(ActionContinuation {
702                    branch_id: branch_id.clone(),
703                    step,
704                    events: outcome.survivors.clone(),
705                    action_ids: action_intents
706                        .iter()
707                        .map(|intent| intent.id.clone())
708                        .collect(),
709                })
710                .map_err(|error| error.to_string())?,
711            )?;
712        }
713        let disposition = match next_at {
714            Some(at) if action_intents.is_empty() => WorkDisposition::Reschedule { at },
715            Some(at) => {
716                Self::set_internal(&mut next_state, "resume_at", json!(at))?;
717                Self::idle_disposition()
718            }
719            None if action_intents.is_empty() && consumed_static_event => WorkDisposition::Complete,
720            None => Self::idle_disposition(),
721        };
722        if !action_intents.is_empty() {
723            let internal = next_state
724                .entry("__workflow")
725                .or_insert_with(|| json!({}))
726                .as_object_mut()
727                .ok_or_else(|| "workflow internal state must be an object".to_string())?;
728            let pending = internal
729                .entry("pending_actions")
730                .or_insert_with(|| json!([]))
731                .as_array_mut()
732                .ok_or_else(|| "workflow pending actions must be an array".to_string())?;
733            pending.extend(action_intents.iter().map(|intent| json!(intent.id)));
734        }
735        let evaluation = EvaluationOutcome {
736            triggered: true,
737            matched: outcome.matched,
738            succeeded: outcome.succeeded && action_intents.is_empty(),
739            action_terminal: resumed_from_actions && action_intents.is_empty(),
740        };
741        Ok(command(
742            item,
743            Value::Object(next_state),
744            action_intents,
745            evaluation,
746            disposition,
747            "workflow.spec_evaluated",
748        ))
749    }
750}
751
752fn command(
753    item: &WorkItem,
754    next_state: Value,
755    action_intents: Vec<ActionIntent>,
756    outcome: EvaluationOutcome,
757    disposition: WorkDisposition,
758    event_type: &str,
759) -> WorkflowTransitionCommand {
760    WorkflowTransitionCommand {
761        consumed_wakeups: if outcome.triggered {
762            item.wakeups
763                .first()
764                .map(|wakeup| wakeup.id.clone())
765                .into_iter()
766                .collect()
767        } else {
768            Vec::new()
769        },
770        delivery_key: format!("spec:{}:{}", item.id, item.state_version),
771        delivery_digest: format!(
772            "{}:{}:{}",
773            item.workflow_revision_digest, item.execution_profile_digest, item.state_version
774        ),
775        event_type: event_type.into(),
776        event_digest: format!("{event_type}:{}:{}", item.id, item.state_version),
777        event_payload: json!({ "spec_id": item.spec_id, "wakeups": item.wakeups.len() }),
778        next_state,
779        action_intents,
780        outcome,
781        disposition,
782    }
783}
784
785#[cfg(test)]
786mod tests {
787    use super::*;
788    use crate::{ControlEpochs, DriverRegistry, PreparedAction, StepNode, StepResult, Wakeup};
789
790    struct PassNode;
791
792    #[async_trait::async_trait]
793    impl StepNode for PassNode {
794        async fn process(&self, event: &Event, _: &crate::WorkflowContext) -> StepResult {
795            StepResult::Pass(event.clone())
796        }
797    }
798
799    struct ActionNode;
800
801    struct FanOutTwo;
802
803    #[async_trait::async_trait]
804    impl StepNode for FanOutTwo {
805        fn produces_fan_out(&self) -> bool {
806            true
807        }
808
809        async fn process(&self, event: &Event, _: &crate::WorkflowContext) -> StepResult {
810            StepResult::FanOut(vec![
811                Event::from_json(json!({"value":event.payload["value"]})),
812                Event::from_json(json!({"value":2})),
813            ])
814        }
815    }
816
817    #[async_trait::async_trait]
818    impl StepNode for ActionNode {
819        async fn process(&self, event: &Event, _: &crate::WorkflowContext) -> StepResult {
820            StepResult::Action {
821                event: event.clone(),
822                action: Box::new(PreparedAction::new(
823                    "execute.demo".parse().unwrap(),
824                    "request-1",
825                    Value::Object(event.payload.clone()),
826                )),
827            }
828        }
829    }
830
831    struct VersionNode(&'static str);
832
833    #[async_trait::async_trait]
834    impl StepNode for VersionNode {
835        async fn process(&self, event: &Event, ctx: &crate::WorkflowContext) -> StepResult {
836            ctx.state
837                .append(&ctx.scoped("seen"), json!(self.0), Some(3), None)
838                .await;
839            StepResult::Pass(event.clone())
840        }
841    }
842
843    fn version_one(_: &Value) -> Result<Box<dyn StepNode>, crate::NodeError> {
844        Ok(Box::new(VersionNode("v1")))
845    }
846
847    fn version_two(_: &Value) -> Result<Box<dyn StepNode>, crate::NodeError> {
848        Ok(Box::new(VersionNode("v2")))
849    }
850
851    fn item(spec_id: &str, config: Value) -> WorkItem {
852        WorkItem {
853            id: "instance".into(),
854            run_id: uuid::Uuid::new_v4().to_string().parse().unwrap(),
855            tenant_id: "tenant".parse().unwrap(),
856            subject_id: "subject".parse().unwrap(),
857            spec_id: spec_id.into(),
858            definition_id: spec_id.into(),
859            workflow_revision: 1,
860            workflow_revision_digest: "digest".into(),
861            execution_profile_id: "profile".into(),
862            execution_profile_revision: 1,
863            execution_profile_digest: "profile-digest".into(),
864            kernel_abi_version: "1".into(),
865            capability_pins: Vec::new(),
866            lifecycle: crate::LifecyclePolicy::run_once(),
867            scheduled_at: chrono::Utc::now(),
868            created_at: chrono::Utc::now(),
869            claimed_at: chrono::Utc::now(),
870            config,
871            state_version: 0,
872            control_epochs: ControlEpochs::default(),
873            cancel_requested: false,
874            lease_version: 1,
875            wakeups: Vec::new(),
876        }
877    }
878
879    fn spec(ingress: &str, ingress_config: Value) -> Spec {
880        Spec::from_json(
881            &json!({
882                "spec_id": "counter", "version": "1",
883                "branches": [{
884                    "branch_id": "__root__",
885                    "nodes": [
886                        {"id": "in", "type": ingress, "config": ingress_config},
887                        {"id": "count", "type": "transform.state_append",
888                         "config": {"key": "seen", "path": "value", "max_len": 3}}
889                    ],
890                    "edges": [{"source": "in", "target": "count"}]
891                }]
892            })
893            .to_string(),
894        )
895        .unwrap()
896    }
897
898    #[tokio::test]
899    async fn published_revisions_select_distinct_graphs_and_reject_conflicting_pins() {
900        let nodes = NodeRegistry::with_builtins();
901        let mut registry = crate::DriverRegistry::<()>::new().with_node_registry(nodes);
902        let first_spec = spec("ingress.event", json!({}));
903        registry.prewarm_spec(&first_spec).unwrap();
904        let first = item("counter", json!({"event": {"value": 7}}));
905        let revision = crate::WorkflowRevision {
906            definition_id: first.definition_id.clone(),
907            revision: 1,
908            content_digest: first.workflow_revision_digest.clone(),
909            kernel_abi_version: "1".into(),
910            dependency_set_digest: "deps".into(),
911            expression_versions: Default::default(),
912            capabilities: Vec::new(),
913            template_provenance: Value::Null,
914            spec: first_spec,
915        };
916        let first_driver = registry
917            .resolve_revision(&first, revision.clone())
918            .await
919            .unwrap();
920        let cached = registry
921            .resolve_revision(&first, revision.clone())
922            .await
923            .unwrap();
924        assert!(Arc::ptr_eq(&first_driver, &cached));
925        let mut second = first.clone();
926        second.workflow_revision = 2;
927        second.workflow_revision_digest = "revision-2".into();
928        let mut second_revision = revision.clone();
929        second_revision.revision = 2;
930        second_revision.content_digest = second.workflow_revision_digest.clone();
931        second_revision.spec.branches[0].nodes[1].config["key"] = json!("second");
932        let second_driver = registry
933            .resolve_revision(&second, second_revision.clone())
934            .await
935            .unwrap();
936        let first_result = first_driver.evaluate(&(), &first).await.unwrap();
937        let second_result = second_driver.evaluate(&(), &second).await.unwrap();
938        assert_eq!(
939            first_result.next_state["state"]["__root__.seen"],
940            json!([7])
941        );
942        assert_eq!(
943            second_result.next_state["state"]["__root__.second"],
944            json!([7])
945        );
946        assert!(second_result.next_state["state"]
947            .get("__root__.seen")
948            .is_none());
949        assert!(registry
950            .resolve_revision(&first, second_revision.clone())
951            .await
952            .is_err());
953        second_revision.spec.branches[0].nodes[1].config["key"] = json!("conflict");
954        assert!(registry
955            .resolve_revision(&second, second_revision.clone())
956            .await
957            .is_err());
958        let mut other_tenant = second.clone();
959        other_tenant.tenant_id = "other-tenant".parse().unwrap();
960        assert!(registry
961            .resolve_revision(&other_tenant, second_revision.clone())
962            .await
963            .is_ok());
964        let mut other_definition = second.clone();
965        other_definition.definition_id = "other-definition".into();
966        second_revision.definition_id = other_definition.definition_id.clone();
967        assert!(registry
968            .resolve_revision(&other_definition, second_revision)
969            .await
970            .is_ok());
971        let mut unsupported = revision.clone();
972        unsupported
973            .expression_versions
974            .insert("transform.state_append".into(), "unknown".into());
975        assert!(registry
976            .resolve_revision(&first, unsupported)
977            .await
978            .is_err());
979        let mut guarded_nodes = NodeRegistry::with_builtins();
980        let mut manifest = CapabilityManifest::action(
981            "transform.state_append",
982            "1",
983            "read-v1",
984            crate::Effect::Read,
985            crate::IdempotencyMode::None,
986            false,
987        );
988        manifest.kind = CapabilityKind::Expression;
989        guarded_nodes.register_capability(manifest).unwrap();
990        let guarded = crate::DriverRegistry::<()>::new().with_node_registry(guarded_nodes);
991        assert!(
992            guarded
993                .resolve_revision(&first, revision.clone())
994                .await
995                .is_err(),
996            "read capabilities need pins even though they do not emit actions"
997        );
998        for version in 3..=259 {
999            let mut next_item = first.clone();
1000            next_item.workflow_revision = version;
1001            next_item.workflow_revision_digest = format!("revision-{version}");
1002            let mut next_revision = revision.clone();
1003            next_revision.revision = version;
1004            next_revision.content_digest = next_item.workflow_revision_digest.clone();
1005            registry
1006                .resolve_revision(&next_item, next_revision)
1007                .await
1008                .unwrap();
1009        }
1010        let (left, right) = tokio::join!(
1011            registry.resolve_revision(&first, revision.clone()),
1012            registry.resolve_revision(&first, revision.clone())
1013        );
1014        let (left, right) = (left.unwrap(), right.unwrap());
1015        assert!(
1016            !Arc::ptr_eq(&left, &first_driver),
1017            "oldest compiled revision is evicted"
1018        );
1019        assert!(
1020            Arc::ptr_eq(&left, &right),
1021            "concurrent resolution shares one compilation"
1022        );
1023        let replaced = registry.with_node_registry(NodeRegistry::empty());
1024        assert!(
1025            replaced
1026                .resolve_revision(&first, revision.clone())
1027                .await
1028                .is_err(),
1029            "changing the node registry must invalidate compiled graphs"
1030        );
1031        let mut unavailable = revision;
1032        unavailable.capabilities.push(crate::CapabilityPin {
1033            id: "missing".into(),
1034            contract_version: "1".into(),
1035            content_digest: "missing-v1".into(),
1036        });
1037        let mut unavailable_item = first;
1038        unavailable_item.capability_pins = unavailable.capabilities.clone();
1039        assert!(replaced
1040            .resolve_revision(&unavailable_item, unavailable)
1041            .await
1042            .is_err());
1043    }
1044
1045    #[tokio::test]
1046    async fn published_revisions_keep_two_versions_of_one_capability_runnable() {
1047        let capability = |version: &str, digest: &str| {
1048            let mut manifest = CapabilityManifest::action(
1049                "transform.state_append",
1050                version,
1051                digest,
1052                crate::Effect::Read,
1053                crate::IdempotencyMode::None,
1054                false,
1055            );
1056            manifest.kind = CapabilityKind::Expression;
1057            manifest
1058        };
1059        let v1 = capability("1", "state-append-v1");
1060        let v2 = capability("2", "state-append-v2");
1061        let pin = |manifest: &CapabilityManifest| crate::CapabilityPin {
1062            id: manifest.id.clone(),
1063            contract_version: manifest.contract_version.clone(),
1064            content_digest: manifest.content_digest.clone(),
1065        };
1066        let mut nodes = NodeRegistry::with_builtins();
1067        let schema = nodes.schema("transform.state_append").cloned().unwrap();
1068        nodes.register_capability(v1.clone()).unwrap();
1069        nodes.register_capability(v2.clone()).unwrap();
1070        nodes
1071            .register_capability_implementation(
1072                pin(&v1),
1073                version_one,
1074                Some(schema.clone()),
1075                None,
1076                false,
1077            )
1078            .unwrap();
1079        nodes
1080            .register_capability_implementation(pin(&v2), version_two, Some(schema), None, false)
1081            .unwrap();
1082        let registry = crate::DriverRegistry::<()>::new().with_node_registry(nodes);
1083        let mut first = item("counter", json!({"event": {"value": 1}}));
1084        first.capability_pins = vec![pin(&v1)];
1085        let mut first_revision = crate::WorkflowRevision {
1086            definition_id: first.definition_id.clone(),
1087            revision: first.workflow_revision,
1088            content_digest: first.workflow_revision_digest.clone(),
1089            kernel_abi_version: "1".into(),
1090            dependency_set_digest: "deps-v1".into(),
1091            expression_versions: Default::default(),
1092            capabilities: first.capability_pins.clone(),
1093            template_provenance: Value::Null,
1094            spec: spec("ingress.event", json!({})),
1095        };
1096        let mut second = first.clone();
1097        second.workflow_revision = 2;
1098        second.workflow_revision_digest = "revision-2".into();
1099        second.capability_pins = vec![pin(&v2)];
1100        let mut second_revision = first_revision.clone();
1101        second_revision.revision = 2;
1102        second_revision.content_digest = second.workflow_revision_digest.clone();
1103        second_revision.dependency_set_digest = "deps-v2".into();
1104        second_revision.capabilities = second.capability_pins.clone();
1105
1106        let first_driver = registry
1107            .resolve_revision(&first, first_revision.clone())
1108            .await
1109            .unwrap();
1110        let second_driver = registry
1111            .resolve_revision(&second, second_revision)
1112            .await
1113            .unwrap();
1114
1115        assert_eq!(
1116            first_driver.evaluate(&(), &first).await.unwrap().next_state["state"]["__root__.seen"],
1117            json!(["v1"])
1118        );
1119        assert_eq!(
1120            second_driver
1121                .evaluate(&(), &second)
1122                .await
1123                .unwrap()
1124                .next_state["state"]["__root__.seen"],
1125            json!(["v2"])
1126        );
1127        first_revision.capabilities = vec![pin(&v2)];
1128        assert!(registry
1129            .resolve_revision(&first, first_revision)
1130            .await
1131            .is_err());
1132    }
1133
1134    #[tokio::test]
1135    async fn old_root_cursor_is_preserved_and_due_schedules_are_not_starved() {
1136        let mut graph = spec("ingress.fixed_rate", json!({"milliseconds": 10_000}));
1137        let mut sibling = spec("ingress.event", json!({})).branches.remove(0);
1138        sibling.branch_id = "events".into();
1139        for node in &mut sibling.nodes {
1140            node.id = format!("events-{}", node.id);
1141        }
1142        for edge in &mut sibling.edges {
1143            edge.source = format!("events-{}", edge.source);
1144            edge.target = format!("events-{}", edge.target);
1145        }
1146        graph.branches.push(sibling);
1147        let driver = SpecDriver::new(&graph, &NodeRegistry::with_builtins()).unwrap();
1148        let start: chrono::DateTime<chrono::Utc> = "2026-01-01T00:00:00Z".parse().unwrap();
1149        let mut work = item("counter", json!({"__workflow":{"catch_up_count":2}}));
1150        work.created_at = start;
1151        work.state_version = 7;
1152        work.scheduled_at = start + chrono::Duration::seconds(100);
1153        work.claimed_at = work.scheduled_at;
1154        let cursors = driver
1155            .schedule_cursors(&work, work.config.as_object().unwrap())
1156            .unwrap();
1157        assert_eq!(cursors["__root__"].next_at, work.scheduled_at);
1158        assert_eq!(cursors["__root__"].catch_up_count, 2);
1159        work.config["__workflow"]["resume_at"] = json!(start + chrono::Duration::seconds(150));
1160        let cursors = driver
1161            .schedule_cursors(&work, work.config.as_object().unwrap())
1162            .unwrap();
1163        assert_eq!(
1164            cursors["__root__"].next_at,
1165            start + chrono::Duration::seconds(150)
1166        );
1167        work.config["__workflow"]
1168            .as_object_mut()
1169            .unwrap()
1170            .remove("resume_at");
1171        work.wakeups.push(serde_json::from_value(json!({"id":"queued", "kind":"delivery", "branch_id":"events", "payload":{"value":7}})).unwrap());
1172        let schedule = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
1173            .await
1174            .unwrap();
1175        assert_eq!(schedule.next_state["last_run"]["branch_id"], "__root__");
1176        assert!(
1177            schedule.consumed_wakeups.is_empty(),
1178            "schedule does not acknowledge queued event"
1179        );
1180        work.config = schedule.next_state;
1181        let event = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
1182            .await
1183            .unwrap();
1184        assert_eq!(event.next_state["last_run"]["branch_id"], "events");
1185        assert_eq!(event.consumed_wakeups, ["queued"]);
1186        work.config = event.next_state;
1187        work.claimed_at += chrono::Duration::seconds(10);
1188        let schedule = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
1189            .await
1190            .unwrap();
1191        assert_eq!(schedule.next_state["last_run"]["branch_id"], "__root__");
1192        assert!(schedule.consumed_wakeups.is_empty());
1193    }
1194
1195    #[test]
1196    fn first_failure_does_not_migrate_a_nonexistent_root_cursor() {
1197        let mut graph = spec("ingress.fixed_rate", json!({"milliseconds":100_000}));
1198        let mut fast = graph.branches[0].clone();
1199        fast.branch_id = "fast".into();
1200        fast.nodes[0].config["milliseconds"] = json!(10_000);
1201        for node in &mut fast.nodes {
1202            node.id = format!("fast-{}", node.id);
1203        }
1204        for edge in &mut fast.edges {
1205            edge.source = format!("fast-{}", edge.source);
1206            edge.target = format!("fast-{}", edge.target);
1207        }
1208        graph.branches.push(fast);
1209        let driver = SpecDriver::new(&graph, &NodeRegistry::with_builtins()).unwrap();
1210        let mut work = item("counter", json!({}));
1211        work.created_at = "2026-01-01T00:00:00Z".parse().unwrap();
1212        work.scheduled_at = work.created_at + chrono::Duration::seconds(10);
1213        work.state_version = 1;
1214        let cursors = driver
1215            .schedule_cursors(&work, work.config.as_object().unwrap())
1216            .unwrap();
1217        assert_eq!(
1218            cursors["__root__"].next_at,
1219            work.created_at + chrono::Duration::seconds(100)
1220        );
1221        assert_eq!(cursors["fast"].next_at, work.scheduled_at);
1222    }
1223
1224    #[tokio::test]
1225    async fn branches_keep_independent_schedule_cursors_across_restart() {
1226        let mut graph = spec("ingress.event", json!({}));
1227        for (branch_id, period) in [("fast", 10_000), ("slow", 25_000)] {
1228            let mut branch = spec("ingress.fixed_rate", json!({"milliseconds": period}))
1229                .branches
1230                .remove(0);
1231            branch.branch_id = branch_id.into();
1232            for node in &mut branch.nodes {
1233                node.id = format!("{branch_id}-{}", node.id);
1234            }
1235            for edge in &mut branch.edges {
1236                edge.source = format!("{branch_id}-{}", edge.source);
1237                edge.target = format!("{branch_id}-{}", edge.target);
1238            }
1239            branch.nodes[1].config["path"] = json!("tick_at");
1240            graph.branches.push(branch);
1241        }
1242        let start: chrono::DateTime<chrono::Utc> = "2026-01-01T00:00:00Z".parse().unwrap();
1243        let mut work = item("counter", json!({}));
1244        work.created_at = start;
1245        work.scheduled_at = start + chrono::Duration::seconds(10);
1246        for (second, branch) in [(10, "fast"), (20, "fast"), (25, "slow"), (30, "fast")] {
1247            work.claimed_at = start + chrono::Duration::seconds(second);
1248            let restarted = SpecDriver::new(&graph, &NodeRegistry::with_builtins()).unwrap();
1249            let command = WorkflowDriver::<()>::evaluate(&restarted, &(), &work)
1250                .await
1251                .unwrap();
1252            assert_eq!(command.next_state["last_run"]["branch_id"], branch);
1253            let WorkDisposition::Reschedule { at } = command.disposition else {
1254                panic!("scheduled branch must retain every cursor");
1255            };
1256            work.config = serde_json::from_str(&command.next_state.to_string()).unwrap();
1257            work.scheduled_at = at;
1258        }
1259        assert_eq!(
1260            work.config["state"]["fast.seen"].as_array().unwrap().len(),
1261            3
1262        );
1263        assert_eq!(
1264            work.config["state"]["slow.seen"],
1265            json!([start + chrono::Duration::seconds(25)])
1266        );
1267        let saved_cursors = work.config["__workflow"]["schedules"].clone();
1268        work.wakeups.push(serde_json::from_value(json!({"id":"event", "kind":"delivery", "branch_id":"__root__", "payload":{"value":99}})).unwrap());
1269        let driver = SpecDriver::new(&graph, &NodeRegistry::with_builtins()).unwrap();
1270        let command = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
1271            .await
1272            .unwrap();
1273        assert_eq!(command.next_state["__workflow"]["schedules"], saved_cursors);
1274        assert_eq!(command.next_state["state"]["__root__.seen"], json!([99]));
1275        work.wakeups.clear();
1276        work.claimed_at = start + chrono::Duration::seconds(31);
1277        let skipped = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
1278            .await
1279            .unwrap();
1280        assert_eq!(skipped.event_type, "workflow.schedule_skipped");
1281        assert_eq!(skipped.next_state["state"], work.config["state"]);
1282    }
1283
1284    #[tokio::test]
1285    async fn delivery_selects_branch_and_rejects_ambiguous_legacy_identity() {
1286        let mut graph = spec("ingress.event", json!({}));
1287        let mut sibling = graph.branches[0].clone();
1288        sibling.branch_id = "other".into();
1289        for node in &mut sibling.nodes {
1290            node.id = format!("other-{}", node.id);
1291        }
1292        for edge in &mut sibling.edges {
1293            edge.source = format!("other-{}", edge.source);
1294            edge.target = format!("other-{}", edge.target);
1295        }
1296        graph.branches.push(sibling);
1297        let driver = SpecDriver::new(&graph, &NodeRegistry::with_builtins()).unwrap();
1298        let mut work = item("counter", json!({}));
1299        work.wakeups.push(
1300            serde_json::from_value(json!({
1301                "id": "delivery", "kind": "delivery", "branch_id": "other",
1302                "payload": {"value": 7}
1303            }))
1304            .unwrap(),
1305        );
1306        let command = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
1307            .await
1308            .unwrap();
1309        assert_eq!(command.next_state["state"]["other.seen"], json!([7]));
1310        assert!(command.next_state["state"]["__root__.seen"].is_null());
1311        work.config = command.next_state;
1312        work.wakeups[0] = serde_json::from_value(json!({
1313            "id": "next", "kind": "delivery", "branch_id": "__root__",
1314            "payload": {"value": 9}
1315        }))
1316        .unwrap();
1317        let restarted = SpecDriver::new(&graph, &NodeRegistry::with_builtins()).unwrap();
1318        let command = WorkflowDriver::<()>::evaluate(&restarted, &(), &work)
1319            .await
1320            .unwrap();
1321        assert_eq!(command.next_state["state"]["other.seen"], json!([7]));
1322        assert_eq!(command.next_state["state"]["__root__.seen"], json!([9]));
1323        for branch in [Value::Null, json!("missing")] {
1324            work.wakeups[0] = serde_json::from_value(json!({
1325                "id": "bad", "kind": "delivery", "branch_id": branch,
1326                "payload": {"value": 1}
1327            }))
1328            .unwrap();
1329            assert!(WorkflowDriver::<()>::evaluate(&driver, &(), &work)
1330                .await
1331                .is_err());
1332        }
1333    }
1334
1335    #[tokio::test]
1336    async fn compiled_driver_rejects_another_spec_identity() {
1337        let registry = NodeRegistry::with_builtins();
1338        let driver = SpecDriver::new(&spec("ingress.event", json!({})), &registry).unwrap();
1339        let result = WorkflowDriver::<()>::evaluate(
1340            &driver,
1341            &(),
1342            &item("different-spec", json!({"event": {"value": 1}})),
1343        )
1344        .await;
1345        assert!(
1346            result.is_err(),
1347            "a compiled graph must not execute another spec"
1348        );
1349    }
1350
1351    #[tokio::test]
1352    async fn cron_spec_reschedules_and_persists_branch_state() {
1353        let registry = NodeRegistry::with_builtins();
1354        let driver = SpecDriver::new(
1355            &spec("ingress.cron", json!({"expression": "0 0 * * * *"})),
1356            &registry,
1357        )
1358        .unwrap();
1359        let first = WorkflowDriver::<()>::evaluate(
1360            &driver,
1361            &(),
1362            &item("counter", json!({"event": {"value": 1}})),
1363        )
1364        .await
1365        .unwrap();
1366        let WorkDisposition::Reschedule { at } = first.disposition else {
1367            panic!("cron spec must reschedule");
1368        };
1369        assert_eq!(first.next_state["last_run"]["terminal"], "completed");
1370        let mut next = item("counter", first.next_state);
1371        next.scheduled_at = at;
1372        next.claimed_at = at;
1373        let second = WorkflowDriver::<()>::evaluate(&driver, &(), &next)
1374            .await
1375            .unwrap();
1376        assert_eq!(
1377            second.next_state["state"]["__root__.seen"],
1378            json!([1]),
1379            "the consumed bootstrap event must not be replayed on a cron tick"
1380        );
1381        assert_eq!(second.next_state["last_run"]["at"], json!(at));
1382    }
1383
1384    #[tokio::test]
1385    async fn instance_config_is_available_and_source_cannot_override_it() {
1386        struct CaptureConfig;
1387        #[async_trait::async_trait]
1388        impl StepNode for CaptureConfig {
1389            async fn process(&self, event: &Event, ctx: &crate::WorkflowContext) -> StepResult {
1390                ctx.state
1391                    .append(
1392                        &ctx.scoped("seen"),
1393                        ctx.config["enabled"].clone(),
1394                        None,
1395                        None,
1396                    )
1397                    .await;
1398                StepResult::Pass(event.clone())
1399            }
1400        }
1401        let mut registry = NodeRegistry::with_builtins();
1402        registry.register_step("transform.capture_config", |_| Ok(Box::new(CaptureConfig)));
1403        let mut graph = spec("ingress.event", json!({}));
1404        graph.branches[0].nodes[1].node_type = "transform.capture_config".into();
1405        graph.branches[0].nodes[1].config = json!({});
1406        let driver = SpecDriver::new(&graph, &registry).unwrap();
1407        for enabled in [true, false] {
1408            let mut work = item("counter", json!({"enabled":enabled}));
1409            work.wakeups.push(Wakeup {
1410                id: "delivery".into(),
1411                branch_id: None,
1412                kind: "delivery".into(),
1413                payload: json!({"enabled":!enabled}),
1414            });
1415            let command = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
1416                .await
1417                .unwrap();
1418            assert_eq!(
1419                command.next_state["state"]["__root__.seen"],
1420                json!([enabled])
1421            );
1422        }
1423    }
1424
1425    #[tokio::test]
1426    async fn event_spec_consumes_a_wakeup_then_waits() {
1427        let registry = NodeRegistry::with_builtins();
1428        let driver = SpecDriver::new(&spec("ingress.event", json!({})), &registry).unwrap();
1429        assert_eq!(
1430            WorkflowDriver::<()>::evaluate(&driver, &(), &item("counter", json!({})))
1431                .await
1432                .unwrap_err(),
1433            "event workflow requires a trigger delivery"
1434        );
1435        let mut work = item("counter", json!({}));
1436        work.wakeups.push(Wakeup {
1437            id: "delivery".into(),
1438            branch_id: None,
1439            kind: "delivery".into(),
1440            payload: json!({"value": 7}),
1441        });
1442        let mut second_wakeup = work.wakeups[0].clone();
1443        second_wakeup.id = "later-delivery".into();
1444        second_wakeup.payload = json!({"value":8});
1445        work.wakeups.push(second_wakeup);
1446        let command = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
1447            .await
1448            .unwrap();
1449        assert_eq!(command.consumed_wakeups, ["delivery"]);
1450        assert_eq!(command.next_state["state"]["__root__.seen"], json!([7]));
1451        assert!(matches!(
1452            command.disposition,
1453            WorkDisposition::Continue { .. }
1454        ));
1455        let mut registry = DriverRegistry::<()>::new();
1456        registry.register(Arc::new(driver)).unwrap();
1457        assert_eq!(registry.spec_ids(), ["counter"]);
1458    }
1459
1460    #[tokio::test]
1461    async fn effectful_steps_are_lifted_into_pinned_intents() {
1462        let mut registry = NodeRegistry::with_builtins();
1463        registry.register_step("guard.auth", |_| Ok(Box::new(PassNode)));
1464        registry.register_side_effect_guard("guard.auth").unwrap();
1465        registry.register_step("execute.demo", |_| Ok(Box::new(ActionNode)));
1466        let manifest = CapabilityManifest::action(
1467            "execute.demo",
1468            "1",
1469            "demo-digest",
1470            crate::Effect::ExternalWrite,
1471            crate::IdempotencyMode::Native,
1472            true,
1473        );
1474        registry.register_capability(manifest.clone()).unwrap();
1475        let effect = Spec::from_json(
1476            &json!({
1477                "spec_id": "effect", "version": "1",
1478                "branches": [{
1479                    "branch_id": "__root__",
1480                        "nodes": [
1481                        {"id": "in", "type": "ingress.event", "config": {}},
1482                        {"id": "auth", "type": "guard.auth", "config": {}},
1483                        {"id": "do", "type": "execute.demo", "config": {}}
1484                    ],
1485                    "edges": [
1486                        {"source": "in", "target": "auth"},
1487                        {"source": "auth", "target": "do"}
1488                    ]
1489                }]
1490            })
1491            .to_string(),
1492        )
1493        .unwrap();
1494        let driver = SpecDriver::new(&effect, &registry).unwrap();
1495        let mut work = item(
1496            "effect",
1497            json!({"objective":"review","event": {"value": 1}}),
1498        );
1499        work.capability_pins.push(crate::CapabilityPin {
1500            id: manifest.id,
1501            contract_version: manifest.contract_version,
1502            content_digest: manifest.content_digest,
1503        });
1504        let command = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
1505            .await
1506            .unwrap();
1507        assert_eq!(command.action_intents.len(), 1);
1508        assert_eq!(command.action_intents[0].state, ActionState::Prepared);
1509        assert_eq!(command.action_intents[0].capability.id, "execute.demo");
1510        assert_eq!(
1511            command.action_intents[0].input,
1512            json!({"objective":"review","value":1})
1513        );
1514        assert_eq!(command.action_intents[0].deadline, None);
1515
1516        let mut waiting = item("effect", command.next_state.clone());
1517        waiting.capability_pins = work.capability_pins.clone();
1518        let waiting_command = WorkflowDriver::<()>::evaluate(&driver, &(), &waiting)
1519            .await
1520            .unwrap();
1521        assert!(waiting_command.action_intents.is_empty());
1522        assert_eq!(waiting_command.event_type, "workflow.action_waiting");
1523
1524        waiting.wakeups.push(Wakeup {
1525            id: "terminal".into(),
1526            branch_id: None,
1527            kind: "timer".into(),
1528            payload: json!({
1529                "kind": "terminal_action",
1530                "observation": {
1531                    "id": "observation",
1532                    "action_intent_id": command.action_intents[0].id,
1533                    "provider_version": "1",
1534                    "observed_at": chrono::Utc::now(),
1535                    "state": "rejected",
1536                    "resource_ref": null,
1537                    "raw_receipt_digest": "receipt",
1538                    "terminal": true,
1539                    "retry_authorized": false
1540                }
1541            }),
1542        });
1543        let mut forged = waiting.clone();
1544        forged.wakeups[0].kind = "delivery".into();
1545        assert!(!forged.is_converging());
1546        let rejected = WorkflowDriver::<()>::evaluate(&driver, &(), &forged)
1547            .await
1548            .unwrap();
1549        assert_eq!(rejected.event_type, "workflow.action_waiting");
1550        assert!(rejected.consumed_wakeups.is_empty());
1551        waiting.wakeups.insert(
1552            0,
1553            Wakeup {
1554                id: "unrelated".into(),
1555                kind: "delivery".into(),
1556                branch_id: None,
1557                payload: json!({"value":999}),
1558            },
1559        );
1560        let terminal = WorkflowDriver::<()>::evaluate(&driver, &(), &waiting)
1561            .await
1562            .unwrap();
1563        assert_eq!(terminal.consumed_wakeups, [waiting.wakeups[1].id.clone()]);
1564        assert!(terminal.action_intents.is_empty());
1565        assert!(matches!(terminal.disposition, WorkDisposition::Complete));
1566        assert!(matches!(
1567            SpecDriver::new(&spec("ingress.cron", json!({})), &registry),
1568            Err(HostError::Schedule { .. })
1569        ));
1570    }
1571
1572    #[tokio::test]
1573    async fn bounded_fanout_actions_resume_with_their_own_results() {
1574        let mut registry = NodeRegistry::with_builtins();
1575        registry.register_step("guard.auth", |_| Ok(Box::new(PassNode)));
1576        registry.register_side_effect_guard("guard.auth").unwrap();
1577        registry.register_step("map.batch", |_| Ok(Box::new(FanOutTwo)));
1578        registry.register_fan_out("map.batch");
1579        registry.register_step("execute.demo", |_| Ok(Box::new(ActionNode)));
1580        let manifest = CapabilityManifest::action(
1581            "execute.demo",
1582            "1",
1583            "demo-digest",
1584            crate::Effect::ExternalWrite,
1585            crate::IdempotencyMode::Native,
1586            true,
1587        );
1588        registry.register_capability(manifest.clone()).unwrap();
1589        let graph = Spec::from_json(&json!({
1590            "spec_id":"fanout-actions","version":"1","branches":[{
1591                "branch_id":"__root__","nodes":[
1592                    {"id":"in","type":"ingress.event","config":{}},
1593                    {"id":"auth","type":"guard.auth","config":{}},
1594                    {"id":"map","type":"map.batch","config":{"count":2}},
1595                    {"id":"action","type":"execute.demo","config":{}},
1596                    {"id":"done","type":"transform.state_append","config":{"key":"results","path":"action_result.result","max_len":2}}
1597                ],"edges":[
1598                    {"source":"in","target":"auth"},{"source":"auth","target":"map"},
1599                    {"source":"map","target":"action"},{"source":"action","target":"done"}
1600                ]
1601            }]
1602        }).to_string()).unwrap();
1603        let driver = SpecDriver::new(&graph, &registry).unwrap();
1604        let mut work = item("fanout-actions", json!({"event":{"value":1}}));
1605        work.capability_pins.push(crate::CapabilityPin {
1606            id: manifest.id,
1607            contract_version: manifest.contract_version,
1608            content_digest: manifest.content_digest,
1609        });
1610        let first = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
1611            .await
1612            .unwrap();
1613        assert_eq!(first.action_intents.len(), 2);
1614        let mut waiting = item("fanout-actions", first.next_state);
1615        waiting.capability_pins = work.capability_pins;
1616        for (index, intent) in first.action_intents.iter().enumerate() {
1617            waiting.wakeups = vec![Wakeup {
1618                id: format!("terminal-{index}"),
1619                branch_id: None,
1620                kind: "timer".into(),
1621                payload: json!({"kind":"terminal_action","observation":{
1622                    "id":format!("observation-{index}"),"action_intent_id":intent.id,
1623                    "provider_version":"1","observed_at":chrono::Utc::now(),
1624                    "state":"succeeded","resource_ref":{"result":format!("result-{index}")},
1625                    "raw_receipt_digest":format!("receipt-{index}"),"terminal":true,
1626                    "retry_authorized":false
1627                }}),
1628            }];
1629            let observed = WorkflowDriver::<()>::evaluate(&driver, &(), &waiting)
1630                .await
1631                .unwrap();
1632            assert!(!observed.outcome.matched);
1633            assert!(!observed.outcome.succeeded);
1634            assert!(!observed.outcome.action_terminal);
1635            waiting.config = observed.next_state;
1636        }
1637        waiting.wakeups.clear();
1638        let resumed = WorkflowDriver::<()>::evaluate(&driver, &(), &waiting)
1639            .await
1640            .unwrap();
1641        assert_eq!(
1642            resumed.next_state["state"]["__root__.results"],
1643            json!(["result-0", "result-1"])
1644        );
1645        assert!(resumed.outcome.matched);
1646        assert!(resumed.outcome.succeeded);
1647        assert!(resumed.outcome.action_terminal);
1648    }
1649
1650    #[tokio::test]
1651    async fn successful_actions_resume_in_order_and_failures_stop_the_chain() {
1652        let mut registry = NodeRegistry::with_builtins();
1653        registry.register_step("guard.auth", |_| Ok(Box::new(PassNode)));
1654        registry.register_side_effect_guard("guard.auth").unwrap();
1655        registry.register_step("execute.demo", |_| Ok(Box::new(ActionNode)));
1656        let manifest = CapabilityManifest::action(
1657            "execute.demo",
1658            "1",
1659            "demo-digest",
1660            crate::Effect::ExternalWrite,
1661            crate::IdempotencyMode::Native,
1662            true,
1663        );
1664        registry.register_capability(manifest.clone()).unwrap();
1665        let graph = Spec::from_json(
1666            &json!({
1667                "spec_id": "chain", "version": "1",
1668                "branches": [{
1669                    "branch_id": "__root__",
1670                    "nodes": [
1671                        {"id":"in","type":"ingress.event","config":{}},
1672                        {"id":"auth","type":"guard.auth","config":{}},
1673                        {"id":"a","type":"execute.demo","config":{}},
1674                        {"id":"b","type":"execute.demo","config":{}},
1675                        {"id":"done","type":"transform.state_append","config":{"key":"results","path":"action_result.result","max_len":3}}
1676                    ],
1677                    "edges": [
1678                        {"source":"in","target":"auth"},
1679                        {"source":"auth","target":"a"},
1680                        {"source":"a","target":"b"},
1681                        {"source":"b","target":"done"}
1682                    ]
1683                }]
1684            }).to_string(),
1685        ).unwrap();
1686        let driver = || SpecDriver::new(&graph, &registry).unwrap();
1687        let pin = crate::CapabilityPin {
1688            id: manifest.id.clone(),
1689            contract_version: manifest.contract_version.clone(),
1690            content_digest: manifest.content_digest.clone(),
1691        };
1692        let mut work = item("chain", json!({"event":{"value":1}}));
1693        work.capability_pins.push(pin.clone());
1694        let first = WorkflowDriver::<()>::evaluate(&driver(), &(), &work)
1695            .await
1696            .unwrap();
1697        assert_eq!(first.action_intents.len(), 1);
1698
1699        let terminal = |action_intent_id: &str, state: &str, result: Value| Wakeup {
1700            id: format!("terminal:{action_intent_id}"),
1701            branch_id: None,
1702            kind: "timer".into(),
1703            payload: json!({"kind":"terminal_action","observation":{
1704                "id":format!("observation:{action_intent_id}"),
1705                "action_intent_id":action_intent_id,
1706                "provider_version":"1","observed_at":chrono::Utc::now(),
1707                "state":state,"resource_ref":result,"raw_receipt_digest":"receipt",
1708                "terminal":true,"retry_authorized":false
1709            }}),
1710        };
1711        let mut observed_a = item("chain", first.next_state);
1712        observed_a.capability_pins.push(pin.clone());
1713        observed_a.state_version = 1;
1714        observed_a.wakeups.push(terminal(
1715            &first.action_intents[0].id,
1716            "succeeded",
1717            json!({"result":"A"}),
1718        ));
1719        let wake_b = WorkflowDriver::<()>::evaluate(&driver(), &(), &observed_a)
1720            .await
1721            .unwrap();
1722        let mut resume_b = item("chain", wake_b.next_state);
1723        resume_b.capability_pins.push(pin.clone());
1724        resume_b.state_version = 2;
1725        let second = WorkflowDriver::<()>::evaluate(&driver(), &(), &resume_b)
1726            .await
1727            .unwrap();
1728        assert_eq!(second.action_intents.len(), 1);
1729        assert_eq!(
1730            second.action_intents[0].input["action_result"]["result"],
1731            "A"
1732        );
1733
1734        let mut observed_b = item("chain", second.next_state.clone());
1735        observed_b.capability_pins.push(pin.clone());
1736        observed_b.state_version = 3;
1737        observed_b.wakeups.push(terminal(
1738            &second.action_intents[0].id,
1739            "succeeded",
1740            json!({"result":"B"}),
1741        ));
1742        let wake_done = WorkflowDriver::<()>::evaluate(&driver(), &(), &observed_b)
1743            .await
1744            .unwrap();
1745        let mut resume_done = item("chain", wake_done.next_state);
1746        resume_done.capability_pins.push(pin.clone());
1747        resume_done.state_version = 4;
1748        let done = WorkflowDriver::<()>::evaluate(&driver(), &(), &resume_done)
1749            .await
1750            .unwrap();
1751        assert_eq!(done.next_state["state"]["__root__.results"], json!(["B"]));
1752
1753        let mut failed = item("chain", second.next_state);
1754        failed.capability_pins.push(pin);
1755        failed.state_version = 3;
1756        failed.wakeups.push(terminal(
1757            &second.action_intents[0].id,
1758            "failed",
1759            Value::Null,
1760        ));
1761        let stopped = WorkflowDriver::<()>::evaluate(&driver(), &(), &failed)
1762            .await
1763            .unwrap();
1764        assert!(matches!(stopped.disposition, WorkDisposition::Complete));
1765        assert!(stopped.next_state["__workflow"]["continuation"].is_null());
1766    }
1767
1768    #[tokio::test]
1769    async fn delivery_does_not_advance_the_cron_cursor() {
1770        let registry = NodeRegistry::with_builtins();
1771        let driver = SpecDriver::new(
1772            &spec("ingress.cron", json!({"expression": "0 0 * * * *"})),
1773            &registry,
1774        )
1775        .unwrap();
1776        let mut work = item("counter", json!({}));
1777        work.scheduled_at = work.claimed_at + chrono::Duration::hours(1);
1778        work.wakeups.push(Wakeup {
1779            id: "delivery".into(),
1780            branch_id: None,
1781            kind: "delivery".into(),
1782            payload: json!({"value": 7}),
1783        });
1784        let command = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
1785            .await
1786            .unwrap();
1787        assert_eq!(command.next_state["state"]["__root__.seen"], json!([7]));
1788        assert!(matches!(
1789            command.disposition,
1790            WorkDisposition::Reschedule { at } if at == work.scheduled_at
1791        ));
1792    }
1793}