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
39impl SpecDriver {
40    /// Compile `spec`; rejects action nodes without an `Action` manifest and funds actions.
41    pub fn new(spec: &Spec, registry: &NodeRegistry) -> Result<Self, HostError> {
42        let host = WorkflowHost::from_spec(spec, registry)?;
43        let schedules = ingress_plans(spec, registry)
44            .into_iter()
45            .map(|plan| {
46                let branch_id = plan.branch_id.clone();
47                (|| {
48                    let required_u64 = |name: &str| {
49                        plan.ingress_config[name]
50                            .as_u64()
51                            .ok_or_else(|| HostError::Schedule {
52                                spec_id: spec.spec_id.clone(),
53                                reason: format!("{} requires integer '{name}'", plan.ingress_type),
54                            })
55                    };
56                    let cadence = match plan.ingress_type.as_str() {
57                        "ingress.cron" => ScheduleCadence::Cron {
58                            expression: plan.ingress_config["expression"]
59                                .as_str()
60                                .ok_or_else(|| HostError::Schedule {
61                                    spec_id: spec.spec_id.clone(),
62                                    reason: "ingress.cron requires a string 'expression'".into(),
63                                })?
64                                .to_owned(),
65                        },
66                        "ingress.fixed_rate" => ScheduleCadence::FixedRate {
67                            milliseconds: required_u64("milliseconds")?,
68                        },
69                        "ingress.fixed_delay" => ScheduleCadence::FixedDelay {
70                            milliseconds: required_u64("milliseconds")?,
71                        },
72                        _ => return Ok(None),
73                    };
74                    let catch_up = match plan.ingress_config["catch_up"].as_str().unwrap_or("once")
75                    {
76                        "skip" => CatchUpPolicy::Skip,
77                        "once" => CatchUpPolicy::CatchUpOnce,
78                        "all" => CatchUpPolicy::CatchUpAll {
79                            limit: plan.ingress_config["catch_up_limit"]
80                                .as_u64()
81                                .and_then(|value| u32::try_from(value).ok())
82                                .ok_or_else(|| HostError::Schedule {
83                                    spec_id: spec.spec_id.clone(),
84                                    reason: "catch_up=all requires positive integer catch_up_limit"
85                                        .into(),
86                                })?,
87                        },
88                        value => {
89                            return Err(HostError::Schedule {
90                                spec_id: spec.spec_id.clone(),
91                                reason: format!("unknown catch_up policy '{value}'"),
92                            })
93                        }
94                    };
95                    let policy = SchedulePolicy {
96                        cadence,
97                        timezone: plan.ingress_config["timezone"]
98                            .as_str()
99                            .unwrap_or("UTC")
100                            .to_owned(),
101                        catch_up,
102                    };
103                    policy.validate().map_err(|error| HostError::Schedule {
104                        spec_id: spec.spec_id.clone(),
105                        reason: error.to_string(),
106                    })?;
107                    Ok::<_, HostError>(Some(policy))
108                })()
109                .map(|policy| policy.map(|policy| (branch_id, policy)))
110            })
111            .collect::<Result<Vec<_>, _>>()?
112            .into_iter()
113            .flatten()
114            .collect();
115        let actions = registry
116            .capability_manifests()
117            .filter(|manifest| manifest.kind == CapabilityKind::Action)
118            .map(|manifest| (manifest.id.clone(), manifest.clone()))
119            .collect();
120        Ok(Self {
121            spec_id: spec.spec_id.clone(),
122            host,
123            schedules,
124            actions,
125        })
126    }
127
128    fn action_intents(
129        &self,
130        item: &WorkItem,
131        requests: Vec<crate::PreparedAction>,
132    ) -> Result<Vec<ActionIntent>, String> {
133        requests
134            .into_iter()
135            .enumerate()
136            .map(|(index, request)| {
137                request.validate().map_err(|error| error.to_string())?;
138                let manifest = self
139                    .actions
140                    .get(request.capability_id.as_str())
141                    .ok_or_else(|| {
142                        format!(
143                            "action step emitted unknown capability '{}'",
144                            request.capability_id
145                        )
146                    })?;
147                let capability = item
148                    .capability_pins
149                    .iter()
150                    .find(|pin| {
151                        pin.id == manifest.id
152                            && pin.contract_version == manifest.contract_version
153                            && pin.content_digest == manifest.content_digest
154                    })
155                    .cloned()
156                    .ok_or_else(|| {
157                        format!(
158                            "workflow revision does not pin action capability '{}'",
159                            manifest.id
160                        )
161                    })?;
162                Ok(ActionIntent {
163                    id: format!("{}:{}:action:{index}", item.id, item.state_version),
164                    tenant_id: item.tenant_id.clone(),
165                    instance_id: item
166                        .id
167                        .parse()
168                        .map_err(|error: af_context::EmptyId| error.to_string())?,
169                    run_id: item.run_id.clone(),
170                    capability,
171                    idempotency_key: format!(
172                        "{}:{}:{}",
173                        item.id, item.state_version, request.idempotency_key
174                    ),
175                    state: ActionState::Prepared,
176                    input: request.input,
177                    effect: manifest.effect,
178                    retry_class: manifest.idempotency_mode,
179                    control_epochs: item.control_epochs,
180                    resource_scope_id: request.resource_scope_id,
181                    lease_epoch: item.lease_version,
182                    action_epoch: item.state_version,
183                    // Provider timeout is per attempt. A product may set a
184                    // separate whole-intent deadline explicitly.
185                    deadline: request.deadline,
186                    reservation: request.reservation,
187                    created_at: item.claimed_at,
188                })
189            })
190            .collect()
191    }
192
193    fn terminal_action(
194        &self,
195        item: &WorkItem,
196        next_state: &mut Map<String, Value>,
197    ) -> Result<Option<WorkflowTransitionCommand>, String> {
198        let Some(wakeup) = item
199            .wakeups
200            .iter()
201            .find(|wakeup| wakeup.kind == "timer" && wakeup.payload["kind"] == "terminal_action")
202        else {
203            return Ok(None);
204        };
205        let observation: crate::ActionObservation =
206            serde_json::from_value(wakeup.payload["observation"].clone())
207                .map_err(|error| format!("terminal action fact: {error}"))?;
208        let Some(internal) = next_state
209            .get_mut("__workflow")
210            .and_then(Value::as_object_mut)
211        else {
212            return Err("terminal action has no durable pending-action state".into());
213        };
214        let pending = internal
215            .get_mut("pending_actions")
216            .and_then(Value::as_array_mut)
217            .ok_or_else(|| "terminal action has no pending action list".to_string())?;
218        let before = pending.len();
219        pending.retain(|id| id.as_str() != Some(&observation.action_intent_id));
220        if pending.len() == before {
221            return Err("terminal action does not belong to this workflow instance".into());
222        }
223        let disposition = if pending.is_empty() {
224            internal
225                .remove("resume_at")
226                .map(serde_json::from_value)
227                .transpose()
228                .map_err(|error| format!("stored workflow resume_at: {error}"))?
229                .map(|at| WorkDisposition::Reschedule { at })
230                .unwrap_or_else(|| {
231                    if internal
232                        .get("static_event_consumed")
233                        .and_then(Value::as_bool)
234                        .unwrap_or(false)
235                    {
236                        WorkDisposition::Complete
237                    } else {
238                        Self::idle_disposition()
239                    }
240                })
241        } else {
242            WorkDisposition::Continue {
243                delay_secs: IDLE_WAIT_SECS,
244            }
245        };
246        let succeeded = matches!(
247            observation.state.as_str(),
248            "completed" | "max_steps_reached" | "succeeded"
249        );
250        let mut selected = item.clone();
251        selected.wakeups = vec![wakeup.clone()];
252        Ok(Some(command(
253            &selected,
254            Value::Object(next_state.clone()),
255            Vec::new(),
256            EvaluationOutcome {
257                triggered: true,
258                matched: true,
259                succeeded,
260                action_terminal: true,
261            },
262            disposition,
263            "workflow.action_observed",
264        )))
265    }
266
267    fn catch_up_count(state: &Map<String, Value>) -> u32 {
268        state
269            .get("__workflow")
270            .and_then(Value::as_object)
271            .and_then(|internal| internal.get("catch_up_count"))
272            .and_then(Value::as_u64)
273            .and_then(|value| u32::try_from(value).ok())
274            .unwrap_or(0)
275    }
276
277    fn set_internal(state: &mut Map<String, Value>, key: &str, value: Value) -> Result<(), String> {
278        let internal = state
279            .entry("__workflow")
280            .or_insert_with(|| json!({}))
281            .as_object_mut()
282            .ok_or_else(|| "workflow internal state must be an object".to_string())?;
283        internal.insert(key.into(), value);
284        Ok(())
285    }
286
287    fn schedule_cursors(
288        &self,
289        item: &WorkItem,
290        state: &Map<String, Value>,
291    ) -> Result<BTreeMap<String, ScheduleCursor>, String> {
292        let saved = state
293            .get("__workflow")
294            .and_then(|value| value.get("schedules"));
295        let mut cursors: BTreeMap<String, ScheduleCursor> = saved
296            .cloned()
297            .map(serde_json::from_value)
298            .transpose()
299            .map_err(|error| format!("stored branch schedule: {error}"))?
300            .unwrap_or_default();
301        for (branch, policy) in &self.schedules {
302            if !cursors.contains_key(branch) {
303                if saved.is_some() {
304                    return Err("stored schedules omit a pinned branch".into());
305                }
306                // Preserve the previously executed root, including a pending action's resume cursor.
307                let legacy_root = self.host.branch_count() == 1
308                    || (saved.is_none()
309                        && item.state_version > 0
310                        && branch == self.host.branch_id()
311                        && (state.contains_key("last_run")
312                            || state.get("__workflow").is_some_and(|internal| {
313                                internal.get("catch_up_count").is_some()
314                                    || internal.get("resume_at").is_some()
315                            })));
316                let next_at = if legacy_root {
317                    state
318                        .get("__workflow")
319                        .and_then(|internal| internal.get("resume_at"))
320                        .cloned()
321                        .map(serde_json::from_value)
322                        .transpose()
323                        .map_err(|error| format!("stored workflow resume_at: {error}"))?
324                        .unwrap_or(item.scheduled_at)
325                } else {
326                    crate::next_scheduled_at(
327                        policy,
328                        item.lifecycle
329                            .starts_at
330                            .map_or(item.created_at, |starts| starts.max(item.created_at)),
331                        None,
332                    )?
333                };
334                cursors.insert(
335                    branch.clone(),
336                    ScheduleCursor {
337                        next_at,
338                        catch_up_count: if legacy_root {
339                            Self::catch_up_count(state)
340                        } else {
341                            0
342                        },
343                    },
344                );
345            }
346        }
347        if cursors
348            .keys()
349            .any(|branch| !self.schedules.contains_key(branch))
350        {
351            return Err("stored schedule refers to an unknown branch".into());
352        }
353        Ok(cursors)
354    }
355
356    fn idle_disposition() -> WorkDisposition {
357        WorkDisposition::Continue {
358            delay_secs: IDLE_WAIT_SECS,
359        }
360    }
361}
362
363#[async_trait]
364impl<Context: Send + Sync> WorkflowDriver<Context> for SpecDriver {
365    fn name(&self) -> &'static str {
366        "spec"
367    }
368
369    fn spec_ids(&self) -> Vec<&str> {
370        vec![&self.spec_id]
371    }
372
373    fn validate_specs(&self) -> Result<(), String> {
374        Ok(())
375    }
376
377    async fn evaluate(
378        &self,
379        _: &Context,
380        item: &WorkItem,
381    ) -> Result<WorkflowTransitionCommand, String> {
382        if item.spec_id != self.spec_id {
383            return Err("compiled spec does not match the claimed spec identity".into());
384        }
385        let mut next_state = match &item.config {
386            Value::Object(map) => map.clone(),
387            Value::Null => Map::new(),
388            _ => return Err("spec instance config must be a JSON object".into()),
389        };
390        if item.cancel_requested {
391            return Ok(command(
392                item,
393                Value::Object(next_state),
394                Vec::new(),
395                EvaluationOutcome::default(),
396                WorkDisposition::Complete,
397                "workflow.cancelled",
398            ));
399        }
400        if let Some(command) = self.terminal_action(item, &mut next_state)? {
401            return Ok(command);
402        }
403        if next_state
404            .get("__workflow")
405            .and_then(Value::as_object)
406            .and_then(|internal| internal.get("pending_actions"))
407            .and_then(Value::as_array)
408            .is_some_and(|pending| !pending.is_empty())
409        {
410            return Ok(command(
411                item,
412                Value::Object(next_state),
413                Vec::new(),
414                EvaluationOutcome::default(),
415                Self::idle_disposition(),
416                "workflow.action_waiting",
417            ));
418        }
419        let mut cursors = self.schedule_cursors(item, &next_state)?;
420        // ponytail: one transition per instance; alternate due schedules and deliveries
421        // under backlog. Use weighted fairness only if consumers need unequal shares.
422        let due = cursors
423            .values()
424            .any(|cursor| cursor.next_at <= item.claimed_at);
425        let prefer_schedule = due
426            && !item.wakeups.is_empty()
427            && next_state
428                .get("__workflow")
429                .and_then(|internal| internal.get("last_dispatch"))
430                .and_then(Value::as_str)
431                != Some("schedule");
432        let mut selected = item.clone();
433        if prefer_schedule {
434            selected.wakeups.clear();
435        }
436        let item = &selected;
437        let earliest = cursors
438            .iter()
439            .min_by_key(|(branch, cursor)| (cursor.next_at, *branch));
440        let branch_id = if let Some(wakeup) = item.wakeups.first() {
441            wakeup
442                .branch_id
443                .as_ref()
444                .map(|id| id.as_str())
445                .or_else(|| (self.host.branch_count() == 1).then(|| self.host.branch_id()))
446                .ok_or_else(|| "multi-branch wakeup requires branch_id".to_string())?
447                .to_owned()
448        } else if let Some((branch, _)) = earliest {
449            branch.clone()
450        } else if self.host.branch_count() == 1 {
451            self.host.branch_id().to_owned()
452        } else {
453            return Err("multi-branch workflow requires a branch-addressed delivery".into());
454        };
455        if !self.host.branch_ids().any(|branch| branch == branch_id) {
456            return Err("wakeup refers to an unknown branch".into());
457        }
458        let schedule = if item.wakeups.is_empty() {
459            if let Some(cursor) = cursors.get_mut(&branch_id) {
460                let decision = if cursor.next_at > item.claimed_at {
461                    crate::ScheduleDecision {
462                        tick_at: None,
463                        next_at: cursor.next_at,
464                        catch_up_count: cursor.catch_up_count,
465                    }
466                } else {
467                    schedule_decision(
468                        &self.schedules[&branch_id],
469                        cursor.next_at,
470                        item.claimed_at,
471                        cursor.catch_up_count,
472                    )?
473                };
474                cursor.next_at = decision.next_at;
475                cursor.catch_up_count = decision.catch_up_count;
476                Some(decision)
477            } else {
478                None
479            }
480        } else {
481            None
482        };
483        Self::set_internal(
484            &mut next_state,
485            "last_dispatch",
486            json!(if item.wakeups.is_empty() {
487                "schedule"
488            } else {
489                "delivery"
490            }),
491        )?;
492        Self::set_internal(&mut next_state, "schedules", json!(cursors))?;
493        let next_at = cursors.values().map(|cursor| cursor.next_at).min();
494        if let Some(decision) = schedule {
495            if decision.tick_at.is_none() {
496                return Ok(command(
497                    item,
498                    Value::Object(next_state),
499                    Vec::new(),
500                    EvaluationOutcome::default(),
501                    WorkDisposition::Reschedule {
502                        at: next_at.expect("selected schedule has a cursor"),
503                    },
504                    "workflow.schedule_skipped",
505                ));
506            }
507        }
508        let state = Arc::new(MemoryState::from_snapshot(
509            next_state
510                .get("state")
511                .and_then(Value::as_object)
512                .cloned()
513                .unwrap_or_default(),
514        ));
515        let wakeup_payload = item.wakeups.first().map(|wakeup| wakeup.payload.clone());
516        let static_event = next_state.get("event").cloned().filter(|_| {
517            !next_state
518                .get("__workflow")
519                .and_then(Value::as_object)
520                .and_then(|internal| internal.get("static_event_consumed"))
521                .and_then(Value::as_bool)
522                .unwrap_or(false)
523        });
524        let consumed_static_event = wakeup_payload.is_none() && static_event.is_some();
525        let payload = wakeup_payload
526            .or(static_event)
527            .or_else(|| {
528                schedule
529                    .and_then(|decision| decision.tick_at)
530                    .map(|tick_at| json!({ "tick_at": tick_at }))
531            })
532            .ok_or_else(|| "event workflow requires a trigger delivery".to_string())?;
533        if consumed_static_event {
534            Self::set_internal(&mut next_state, "static_event_consumed", json!(true))?;
535        }
536        let context = self.host.context_for(&branch_id, state.clone());
537        let outcome = self
538            .host
539            .run_event(&context, Event::from_json(payload))
540            .await
541            .ok_or_else(|| "spec has no selected branch".to_string())?;
542        let (terminal, exit_reason) = match &outcome.terminal {
543            Terminal::Completed => ("completed", Value::Null),
544            Terminal::Dropped { node_id, reason } => {
545                ("dropped", json!({ "node_id": node_id, "reason": reason }))
546            }
547        };
548        next_state.insert("state".into(), Value::Object(state.snapshot()));
549        next_state.insert(
550            "last_run".into(),
551            json!({
552                "branch_id": branch_id,
553                "terminal": terminal,
554                "exit": exit_reason,
555                "steps_run": outcome.steps_run,
556                "survivors": outcome.survivors.len(),
557                "at": item.claimed_at,
558            }),
559        );
560        let action_intents = self.action_intents(item, outcome.actions)?;
561        let disposition = match next_at {
562            Some(at) if action_intents.is_empty() => WorkDisposition::Reschedule { at },
563            Some(at) => {
564                Self::set_internal(&mut next_state, "resume_at", json!(at))?;
565                Self::idle_disposition()
566            }
567            None if action_intents.is_empty() && consumed_static_event => WorkDisposition::Complete,
568            None => Self::idle_disposition(),
569        };
570        if !action_intents.is_empty() {
571            let internal = next_state
572                .entry("__workflow")
573                .or_insert_with(|| json!({}))
574                .as_object_mut()
575                .ok_or_else(|| "workflow internal state must be an object".to_string())?;
576            let pending = internal
577                .entry("pending_actions")
578                .or_insert_with(|| json!([]))
579                .as_array_mut()
580                .ok_or_else(|| "workflow pending actions must be an array".to_string())?;
581            pending.extend(action_intents.iter().map(|intent| json!(intent.id)));
582        }
583        let evaluation = EvaluationOutcome {
584            triggered: true,
585            matched: outcome.matched,
586            succeeded: outcome.succeeded && action_intents.is_empty(),
587            action_terminal: false,
588        };
589        Ok(command(
590            item,
591            Value::Object(next_state),
592            action_intents,
593            evaluation,
594            disposition,
595            "workflow.spec_evaluated",
596        ))
597    }
598}
599
600fn command(
601    item: &WorkItem,
602    next_state: Value,
603    action_intents: Vec<ActionIntent>,
604    outcome: EvaluationOutcome,
605    disposition: WorkDisposition,
606    event_type: &str,
607) -> WorkflowTransitionCommand {
608    WorkflowTransitionCommand {
609        consumed_wakeups: if outcome.triggered {
610            item.wakeups
611                .first()
612                .map(|wakeup| wakeup.id.clone())
613                .into_iter()
614                .collect()
615        } else {
616            Vec::new()
617        },
618        delivery_key: format!("spec:{}:{}", item.id, item.state_version),
619        delivery_digest: format!(
620            "{}:{}:{}",
621            item.workflow_revision_digest, item.execution_profile_digest, item.state_version
622        ),
623        event_type: event_type.into(),
624        event_digest: format!("{event_type}:{}:{}", item.id, item.state_version),
625        event_payload: json!({ "spec_id": item.spec_id, "wakeups": item.wakeups.len() }),
626        next_state,
627        action_intents,
628        outcome,
629        disposition,
630    }
631}
632
633#[cfg(test)]
634mod tests {
635    use super::*;
636    use crate::{ControlEpochs, DriverRegistry, PreparedAction, StepNode, StepResult, Wakeup};
637
638    struct PassNode;
639
640    #[async_trait::async_trait]
641    impl StepNode for PassNode {
642        async fn process(&self, event: &Event, _: &crate::WorkflowContext) -> StepResult {
643            StepResult::Pass(event.clone())
644        }
645    }
646
647    struct ActionNode;
648
649    #[async_trait::async_trait]
650    impl StepNode for ActionNode {
651        async fn process(&self, event: &Event, _: &crate::WorkflowContext) -> StepResult {
652            StepResult::Action {
653                event: event.clone(),
654                action: Box::new(PreparedAction::new(
655                    "execute.demo".parse().unwrap(),
656                    "request-1",
657                    json!({"value": 1}),
658                )),
659            }
660        }
661    }
662
663    struct VersionNode(&'static str);
664
665    #[async_trait::async_trait]
666    impl StepNode for VersionNode {
667        async fn process(&self, event: &Event, ctx: &crate::WorkflowContext) -> StepResult {
668            ctx.state
669                .append(&ctx.scoped("seen"), json!(self.0), Some(3), None)
670                .await;
671            StepResult::Pass(event.clone())
672        }
673    }
674
675    fn version_one(_: &Value) -> Result<Box<dyn StepNode>, crate::NodeError> {
676        Ok(Box::new(VersionNode("v1")))
677    }
678
679    fn version_two(_: &Value) -> Result<Box<dyn StepNode>, crate::NodeError> {
680        Ok(Box::new(VersionNode("v2")))
681    }
682
683    fn item(spec_id: &str, config: Value) -> WorkItem {
684        WorkItem {
685            id: "instance".into(),
686            run_id: uuid::Uuid::new_v4().to_string().parse().unwrap(),
687            tenant_id: "tenant".parse().unwrap(),
688            subject_id: "subject".parse().unwrap(),
689            spec_id: spec_id.into(),
690            definition_id: spec_id.into(),
691            workflow_revision: 1,
692            workflow_revision_digest: "digest".into(),
693            execution_profile_id: "profile".into(),
694            execution_profile_revision: 1,
695            execution_profile_digest: "profile-digest".into(),
696            kernel_abi_version: "1".into(),
697            capability_pins: Vec::new(),
698            lifecycle: crate::LifecyclePolicy::run_once(),
699            scheduled_at: chrono::Utc::now(),
700            created_at: chrono::Utc::now(),
701            claimed_at: chrono::Utc::now(),
702            config,
703            state_version: 0,
704            control_epochs: ControlEpochs::default(),
705            cancel_requested: false,
706            lease_version: 1,
707            wakeups: Vec::new(),
708        }
709    }
710
711    fn spec(ingress: &str, ingress_config: Value) -> Spec {
712        Spec::from_json(
713            &json!({
714                "spec_id": "counter", "version": "1",
715                "branches": [{
716                    "branch_id": "__root__",
717                    "nodes": [
718                        {"id": "in", "type": ingress, "config": ingress_config},
719                        {"id": "count", "type": "transform.state_append",
720                         "config": {"key": "seen", "path": "value", "max_len": 3}}
721                    ],
722                    "edges": [{"source": "in", "target": "count"}]
723                }]
724            })
725            .to_string(),
726        )
727        .unwrap()
728    }
729
730    #[tokio::test]
731    async fn published_revisions_select_distinct_graphs_and_reject_conflicting_pins() {
732        let nodes = NodeRegistry::with_builtins();
733        let mut registry = crate::DriverRegistry::<()>::new().with_node_registry(nodes);
734        let first_spec = spec("ingress.event", json!({}));
735        registry.prewarm_spec(&first_spec).unwrap();
736        let first = item("counter", json!({"event": {"value": 7}}));
737        let revision = crate::WorkflowRevision {
738            definition_id: first.definition_id.clone(),
739            revision: 1,
740            content_digest: first.workflow_revision_digest.clone(),
741            kernel_abi_version: "1".into(),
742            dependency_set_digest: "deps".into(),
743            expression_versions: Default::default(),
744            capabilities: Vec::new(),
745            template_provenance: Value::Null,
746            spec: first_spec,
747        };
748        let first_driver = registry
749            .resolve_revision(&first, revision.clone())
750            .await
751            .unwrap();
752        let cached = registry
753            .resolve_revision(&first, revision.clone())
754            .await
755            .unwrap();
756        assert!(Arc::ptr_eq(&first_driver, &cached));
757        let mut second = first.clone();
758        second.workflow_revision = 2;
759        second.workflow_revision_digest = "revision-2".into();
760        let mut second_revision = revision.clone();
761        second_revision.revision = 2;
762        second_revision.content_digest = second.workflow_revision_digest.clone();
763        second_revision.spec.branches[0].nodes[1].config["key"] = json!("second");
764        let second_driver = registry
765            .resolve_revision(&second, second_revision.clone())
766            .await
767            .unwrap();
768        let first_result = first_driver.evaluate(&(), &first).await.unwrap();
769        let second_result = second_driver.evaluate(&(), &second).await.unwrap();
770        assert_eq!(
771            first_result.next_state["state"]["__root__.seen"],
772            json!([7])
773        );
774        assert_eq!(
775            second_result.next_state["state"]["__root__.second"],
776            json!([7])
777        );
778        assert!(second_result.next_state["state"]
779            .get("__root__.seen")
780            .is_none());
781        assert!(registry
782            .resolve_revision(&first, second_revision.clone())
783            .await
784            .is_err());
785        second_revision.spec.branches[0].nodes[1].config["key"] = json!("conflict");
786        assert!(registry
787            .resolve_revision(&second, second_revision.clone())
788            .await
789            .is_err());
790        let mut other_tenant = second.clone();
791        other_tenant.tenant_id = "other-tenant".parse().unwrap();
792        assert!(registry
793            .resolve_revision(&other_tenant, second_revision.clone())
794            .await
795            .is_ok());
796        let mut other_definition = second.clone();
797        other_definition.definition_id = "other-definition".into();
798        second_revision.definition_id = other_definition.definition_id.clone();
799        assert!(registry
800            .resolve_revision(&other_definition, second_revision)
801            .await
802            .is_ok());
803        let mut unsupported = revision.clone();
804        unsupported
805            .expression_versions
806            .insert("transform.state_append".into(), "unknown".into());
807        assert!(registry
808            .resolve_revision(&first, unsupported)
809            .await
810            .is_err());
811        let mut guarded_nodes = NodeRegistry::with_builtins();
812        let mut manifest = CapabilityManifest::action(
813            "transform.state_append",
814            "1",
815            "read-v1",
816            crate::Effect::Read,
817            crate::IdempotencyMode::None,
818            false,
819        );
820        manifest.kind = CapabilityKind::Expression;
821        guarded_nodes.register_capability(manifest).unwrap();
822        let guarded = crate::DriverRegistry::<()>::new().with_node_registry(guarded_nodes);
823        assert!(
824            guarded
825                .resolve_revision(&first, revision.clone())
826                .await
827                .is_err(),
828            "read capabilities need pins even though they do not emit actions"
829        );
830        for version in 3..=259 {
831            let mut next_item = first.clone();
832            next_item.workflow_revision = version;
833            next_item.workflow_revision_digest = format!("revision-{version}");
834            let mut next_revision = revision.clone();
835            next_revision.revision = version;
836            next_revision.content_digest = next_item.workflow_revision_digest.clone();
837            registry
838                .resolve_revision(&next_item, next_revision)
839                .await
840                .unwrap();
841        }
842        let (left, right) = tokio::join!(
843            registry.resolve_revision(&first, revision.clone()),
844            registry.resolve_revision(&first, revision.clone())
845        );
846        let (left, right) = (left.unwrap(), right.unwrap());
847        assert!(
848            !Arc::ptr_eq(&left, &first_driver),
849            "oldest compiled revision is evicted"
850        );
851        assert!(
852            Arc::ptr_eq(&left, &right),
853            "concurrent resolution shares one compilation"
854        );
855        let replaced = registry.with_node_registry(NodeRegistry::empty());
856        assert!(
857            replaced
858                .resolve_revision(&first, revision.clone())
859                .await
860                .is_err(),
861            "changing the node registry must invalidate compiled graphs"
862        );
863        let mut unavailable = revision;
864        unavailable.capabilities.push(crate::CapabilityPin {
865            id: "missing".into(),
866            contract_version: "1".into(),
867            content_digest: "missing-v1".into(),
868        });
869        let mut unavailable_item = first;
870        unavailable_item.capability_pins = unavailable.capabilities.clone();
871        assert!(replaced
872            .resolve_revision(&unavailable_item, unavailable)
873            .await
874            .is_err());
875    }
876
877    #[tokio::test]
878    async fn published_revisions_keep_two_versions_of_one_capability_runnable() {
879        let capability = |version: &str, digest: &str| {
880            let mut manifest = CapabilityManifest::action(
881                "transform.state_append",
882                version,
883                digest,
884                crate::Effect::Read,
885                crate::IdempotencyMode::None,
886                false,
887            );
888            manifest.kind = CapabilityKind::Expression;
889            manifest
890        };
891        let v1 = capability("1", "state-append-v1");
892        let v2 = capability("2", "state-append-v2");
893        let pin = |manifest: &CapabilityManifest| crate::CapabilityPin {
894            id: manifest.id.clone(),
895            contract_version: manifest.contract_version.clone(),
896            content_digest: manifest.content_digest.clone(),
897        };
898        let mut nodes = NodeRegistry::with_builtins();
899        let schema = nodes.schema("transform.state_append").cloned().unwrap();
900        nodes.register_capability(v1.clone()).unwrap();
901        nodes.register_capability(v2.clone()).unwrap();
902        nodes
903            .register_capability_implementation(
904                pin(&v1),
905                version_one,
906                Some(schema.clone()),
907                None,
908                false,
909            )
910            .unwrap();
911        nodes
912            .register_capability_implementation(pin(&v2), version_two, Some(schema), None, false)
913            .unwrap();
914        let registry = crate::DriverRegistry::<()>::new().with_node_registry(nodes);
915        let mut first = item("counter", json!({"event": {"value": 1}}));
916        first.capability_pins = vec![pin(&v1)];
917        let mut first_revision = crate::WorkflowRevision {
918            definition_id: first.definition_id.clone(),
919            revision: first.workflow_revision,
920            content_digest: first.workflow_revision_digest.clone(),
921            kernel_abi_version: "1".into(),
922            dependency_set_digest: "deps-v1".into(),
923            expression_versions: Default::default(),
924            capabilities: first.capability_pins.clone(),
925            template_provenance: Value::Null,
926            spec: spec("ingress.event", json!({})),
927        };
928        let mut second = first.clone();
929        second.workflow_revision = 2;
930        second.workflow_revision_digest = "revision-2".into();
931        second.capability_pins = vec![pin(&v2)];
932        let mut second_revision = first_revision.clone();
933        second_revision.revision = 2;
934        second_revision.content_digest = second.workflow_revision_digest.clone();
935        second_revision.dependency_set_digest = "deps-v2".into();
936        second_revision.capabilities = second.capability_pins.clone();
937
938        let first_driver = registry
939            .resolve_revision(&first, first_revision.clone())
940            .await
941            .unwrap();
942        let second_driver = registry
943            .resolve_revision(&second, second_revision)
944            .await
945            .unwrap();
946
947        assert_eq!(
948            first_driver.evaluate(&(), &first).await.unwrap().next_state["state"]["__root__.seen"],
949            json!(["v1"])
950        );
951        assert_eq!(
952            second_driver
953                .evaluate(&(), &second)
954                .await
955                .unwrap()
956                .next_state["state"]["__root__.seen"],
957            json!(["v2"])
958        );
959        first_revision.capabilities = vec![pin(&v2)];
960        assert!(registry
961            .resolve_revision(&first, first_revision)
962            .await
963            .is_err());
964    }
965
966    #[tokio::test]
967    async fn old_root_cursor_is_preserved_and_due_schedules_are_not_starved() {
968        let mut graph = spec("ingress.fixed_rate", json!({"milliseconds": 10_000}));
969        let mut sibling = spec("ingress.event", json!({})).branches.remove(0);
970        sibling.branch_id = "events".into();
971        for node in &mut sibling.nodes {
972            node.id = format!("events-{}", node.id);
973        }
974        for edge in &mut sibling.edges {
975            edge.source = format!("events-{}", edge.source);
976            edge.target = format!("events-{}", edge.target);
977        }
978        graph.branches.push(sibling);
979        let driver = SpecDriver::new(&graph, &NodeRegistry::with_builtins()).unwrap();
980        let start: chrono::DateTime<chrono::Utc> = "2026-01-01T00:00:00Z".parse().unwrap();
981        let mut work = item("counter", json!({"__workflow":{"catch_up_count":2}}));
982        work.created_at = start;
983        work.state_version = 7;
984        work.scheduled_at = start + chrono::Duration::seconds(100);
985        work.claimed_at = work.scheduled_at;
986        let cursors = driver
987            .schedule_cursors(&work, work.config.as_object().unwrap())
988            .unwrap();
989        assert_eq!(cursors["__root__"].next_at, work.scheduled_at);
990        assert_eq!(cursors["__root__"].catch_up_count, 2);
991        work.config["__workflow"]["resume_at"] = json!(start + chrono::Duration::seconds(150));
992        let cursors = driver
993            .schedule_cursors(&work, work.config.as_object().unwrap())
994            .unwrap();
995        assert_eq!(
996            cursors["__root__"].next_at,
997            start + chrono::Duration::seconds(150)
998        );
999        work.config["__workflow"]
1000            .as_object_mut()
1001            .unwrap()
1002            .remove("resume_at");
1003        work.wakeups.push(serde_json::from_value(json!({"id":"queued", "kind":"delivery", "branch_id":"events", "payload":{"value":7}})).unwrap());
1004        let schedule = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
1005            .await
1006            .unwrap();
1007        assert_eq!(schedule.next_state["last_run"]["branch_id"], "__root__");
1008        assert!(
1009            schedule.consumed_wakeups.is_empty(),
1010            "schedule does not acknowledge queued event"
1011        );
1012        work.config = schedule.next_state;
1013        let event = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
1014            .await
1015            .unwrap();
1016        assert_eq!(event.next_state["last_run"]["branch_id"], "events");
1017        assert_eq!(event.consumed_wakeups, ["queued"]);
1018        work.config = event.next_state;
1019        work.claimed_at += chrono::Duration::seconds(10);
1020        let schedule = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
1021            .await
1022            .unwrap();
1023        assert_eq!(schedule.next_state["last_run"]["branch_id"], "__root__");
1024        assert!(schedule.consumed_wakeups.is_empty());
1025    }
1026
1027    #[test]
1028    fn first_failure_does_not_migrate_a_nonexistent_root_cursor() {
1029        let mut graph = spec("ingress.fixed_rate", json!({"milliseconds":100_000}));
1030        let mut fast = graph.branches[0].clone();
1031        fast.branch_id = "fast".into();
1032        fast.nodes[0].config["milliseconds"] = json!(10_000);
1033        for node in &mut fast.nodes {
1034            node.id = format!("fast-{}", node.id);
1035        }
1036        for edge in &mut fast.edges {
1037            edge.source = format!("fast-{}", edge.source);
1038            edge.target = format!("fast-{}", edge.target);
1039        }
1040        graph.branches.push(fast);
1041        let driver = SpecDriver::new(&graph, &NodeRegistry::with_builtins()).unwrap();
1042        let mut work = item("counter", json!({}));
1043        work.created_at = "2026-01-01T00:00:00Z".parse().unwrap();
1044        work.scheduled_at = work.created_at + chrono::Duration::seconds(10);
1045        work.state_version = 1;
1046        let cursors = driver
1047            .schedule_cursors(&work, work.config.as_object().unwrap())
1048            .unwrap();
1049        assert_eq!(
1050            cursors["__root__"].next_at,
1051            work.created_at + chrono::Duration::seconds(100)
1052        );
1053        assert_eq!(cursors["fast"].next_at, work.scheduled_at);
1054    }
1055
1056    #[tokio::test]
1057    async fn branches_keep_independent_schedule_cursors_across_restart() {
1058        let mut graph = spec("ingress.event", json!({}));
1059        for (branch_id, period) in [("fast", 10_000), ("slow", 25_000)] {
1060            let mut branch = spec("ingress.fixed_rate", json!({"milliseconds": period}))
1061                .branches
1062                .remove(0);
1063            branch.branch_id = branch_id.into();
1064            for node in &mut branch.nodes {
1065                node.id = format!("{branch_id}-{}", node.id);
1066            }
1067            for edge in &mut branch.edges {
1068                edge.source = format!("{branch_id}-{}", edge.source);
1069                edge.target = format!("{branch_id}-{}", edge.target);
1070            }
1071            branch.nodes[1].config["path"] = json!("tick_at");
1072            graph.branches.push(branch);
1073        }
1074        let start: chrono::DateTime<chrono::Utc> = "2026-01-01T00:00:00Z".parse().unwrap();
1075        let mut work = item("counter", json!({}));
1076        work.created_at = start;
1077        work.scheduled_at = start + chrono::Duration::seconds(10);
1078        for (second, branch) in [(10, "fast"), (20, "fast"), (25, "slow"), (30, "fast")] {
1079            work.claimed_at = start + chrono::Duration::seconds(second);
1080            let restarted = SpecDriver::new(&graph, &NodeRegistry::with_builtins()).unwrap();
1081            let command = WorkflowDriver::<()>::evaluate(&restarted, &(), &work)
1082                .await
1083                .unwrap();
1084            assert_eq!(command.next_state["last_run"]["branch_id"], branch);
1085            let WorkDisposition::Reschedule { at } = command.disposition else {
1086                panic!("scheduled branch must retain every cursor");
1087            };
1088            work.config = serde_json::from_str(&command.next_state.to_string()).unwrap();
1089            work.scheduled_at = at;
1090        }
1091        assert_eq!(
1092            work.config["state"]["fast.seen"].as_array().unwrap().len(),
1093            3
1094        );
1095        assert_eq!(
1096            work.config["state"]["slow.seen"],
1097            json!([start + chrono::Duration::seconds(25)])
1098        );
1099        let saved_cursors = work.config["__workflow"]["schedules"].clone();
1100        work.wakeups.push(serde_json::from_value(json!({"id":"event", "kind":"delivery", "branch_id":"__root__", "payload":{"value":99}})).unwrap());
1101        let driver = SpecDriver::new(&graph, &NodeRegistry::with_builtins()).unwrap();
1102        let command = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
1103            .await
1104            .unwrap();
1105        assert_eq!(command.next_state["__workflow"]["schedules"], saved_cursors);
1106        assert_eq!(command.next_state["state"]["__root__.seen"], json!([99]));
1107        work.wakeups.clear();
1108        work.claimed_at = start + chrono::Duration::seconds(31);
1109        let skipped = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
1110            .await
1111            .unwrap();
1112        assert_eq!(skipped.event_type, "workflow.schedule_skipped");
1113        assert_eq!(skipped.next_state["state"], work.config["state"]);
1114    }
1115
1116    #[tokio::test]
1117    async fn delivery_selects_branch_and_rejects_ambiguous_legacy_identity() {
1118        let mut graph = spec("ingress.event", json!({}));
1119        let mut sibling = graph.branches[0].clone();
1120        sibling.branch_id = "other".into();
1121        for node in &mut sibling.nodes {
1122            node.id = format!("other-{}", node.id);
1123        }
1124        for edge in &mut sibling.edges {
1125            edge.source = format!("other-{}", edge.source);
1126            edge.target = format!("other-{}", edge.target);
1127        }
1128        graph.branches.push(sibling);
1129        let driver = SpecDriver::new(&graph, &NodeRegistry::with_builtins()).unwrap();
1130        let mut work = item("counter", json!({}));
1131        work.wakeups.push(
1132            serde_json::from_value(json!({
1133                "id": "delivery", "kind": "delivery", "branch_id": "other",
1134                "payload": {"value": 7}
1135            }))
1136            .unwrap(),
1137        );
1138        let command = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
1139            .await
1140            .unwrap();
1141        assert_eq!(command.next_state["state"]["other.seen"], json!([7]));
1142        assert!(command.next_state["state"]["__root__.seen"].is_null());
1143        work.config = command.next_state;
1144        work.wakeups[0] = serde_json::from_value(json!({
1145            "id": "next", "kind": "delivery", "branch_id": "__root__",
1146            "payload": {"value": 9}
1147        }))
1148        .unwrap();
1149        let restarted = SpecDriver::new(&graph, &NodeRegistry::with_builtins()).unwrap();
1150        let command = WorkflowDriver::<()>::evaluate(&restarted, &(), &work)
1151            .await
1152            .unwrap();
1153        assert_eq!(command.next_state["state"]["other.seen"], json!([7]));
1154        assert_eq!(command.next_state["state"]["__root__.seen"], json!([9]));
1155        for branch in [Value::Null, json!("missing")] {
1156            work.wakeups[0] = serde_json::from_value(json!({
1157                "id": "bad", "kind": "delivery", "branch_id": branch,
1158                "payload": {"value": 1}
1159            }))
1160            .unwrap();
1161            assert!(WorkflowDriver::<()>::evaluate(&driver, &(), &work)
1162                .await
1163                .is_err());
1164        }
1165    }
1166
1167    #[tokio::test]
1168    async fn compiled_driver_rejects_another_spec_identity() {
1169        let registry = NodeRegistry::with_builtins();
1170        let driver = SpecDriver::new(&spec("ingress.event", json!({})), &registry).unwrap();
1171        let result = WorkflowDriver::<()>::evaluate(
1172            &driver,
1173            &(),
1174            &item("different-spec", json!({"event": {"value": 1}})),
1175        )
1176        .await;
1177        assert!(
1178            result.is_err(),
1179            "a compiled graph must not execute another spec"
1180        );
1181    }
1182
1183    #[tokio::test]
1184    async fn cron_spec_reschedules_and_persists_branch_state() {
1185        let registry = NodeRegistry::with_builtins();
1186        let driver = SpecDriver::new(
1187            &spec("ingress.cron", json!({"expression": "0 0 * * * *"})),
1188            &registry,
1189        )
1190        .unwrap();
1191        let first = WorkflowDriver::<()>::evaluate(
1192            &driver,
1193            &(),
1194            &item("counter", json!({"event": {"value": 1}})),
1195        )
1196        .await
1197        .unwrap();
1198        let WorkDisposition::Reschedule { at } = first.disposition else {
1199            panic!("cron spec must reschedule");
1200        };
1201        assert_eq!(first.next_state["last_run"]["terminal"], "completed");
1202        let mut next = item("counter", first.next_state);
1203        next.scheduled_at = at;
1204        next.claimed_at = at;
1205        let second = WorkflowDriver::<()>::evaluate(&driver, &(), &next)
1206            .await
1207            .unwrap();
1208        assert_eq!(
1209            second.next_state["state"]["__root__.seen"],
1210            json!([1]),
1211            "the consumed bootstrap event must not be replayed on a cron tick"
1212        );
1213        assert_eq!(second.next_state["last_run"]["at"], json!(at));
1214    }
1215
1216    #[tokio::test]
1217    async fn event_spec_consumes_a_wakeup_then_waits() {
1218        let registry = NodeRegistry::with_builtins();
1219        let driver = SpecDriver::new(&spec("ingress.event", json!({})), &registry).unwrap();
1220        assert_eq!(
1221            WorkflowDriver::<()>::evaluate(&driver, &(), &item("counter", json!({})))
1222                .await
1223                .unwrap_err(),
1224            "event workflow requires a trigger delivery"
1225        );
1226        let mut work = item("counter", json!({}));
1227        work.wakeups.push(Wakeup {
1228            id: "delivery".into(),
1229            branch_id: None,
1230            kind: "delivery".into(),
1231            payload: json!({"value": 7}),
1232        });
1233        let mut second_wakeup = work.wakeups[0].clone();
1234        second_wakeup.id = "later-delivery".into();
1235        second_wakeup.payload = json!({"value":8});
1236        work.wakeups.push(second_wakeup);
1237        let command = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
1238            .await
1239            .unwrap();
1240        assert_eq!(command.consumed_wakeups, ["delivery"]);
1241        assert_eq!(command.next_state["state"]["__root__.seen"], json!([7]));
1242        assert!(matches!(
1243            command.disposition,
1244            WorkDisposition::Continue { .. }
1245        ));
1246        let mut registry = DriverRegistry::<()>::new();
1247        registry.register(Arc::new(driver)).unwrap();
1248        assert_eq!(registry.spec_ids(), ["counter"]);
1249    }
1250
1251    #[tokio::test]
1252    async fn effectful_steps_are_lifted_into_pinned_intents() {
1253        let mut registry = NodeRegistry::with_builtins();
1254        registry.register_step("guard.auth", |_| Ok(Box::new(PassNode)));
1255        registry.register_side_effect_guard("guard.auth");
1256        registry.register_step("execute.demo", |_| Ok(Box::new(ActionNode)));
1257        let manifest = CapabilityManifest::action(
1258            "execute.demo",
1259            "1",
1260            "demo-digest",
1261            crate::Effect::ExternalWrite,
1262            crate::IdempotencyMode::Native,
1263            true,
1264        );
1265        registry.register_capability(manifest.clone()).unwrap();
1266        let effect = Spec::from_json(
1267            &json!({
1268                "spec_id": "effect", "version": "1",
1269                "branches": [{
1270                    "branch_id": "__root__",
1271                        "nodes": [
1272                        {"id": "in", "type": "ingress.event", "config": {}},
1273                        {"id": "auth", "type": "guard.auth", "config": {}},
1274                        {"id": "do", "type": "execute.demo", "config": {}}
1275                    ],
1276                    "edges": [
1277                        {"source": "in", "target": "auth"},
1278                        {"source": "auth", "target": "do"}
1279                    ]
1280                }]
1281            })
1282            .to_string(),
1283        )
1284        .unwrap();
1285        let driver = SpecDriver::new(&effect, &registry).unwrap();
1286        let mut work = item("effect", json!({"event": {"value": 1}}));
1287        work.capability_pins.push(crate::CapabilityPin {
1288            id: manifest.id,
1289            contract_version: manifest.contract_version,
1290            content_digest: manifest.content_digest,
1291        });
1292        let command = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
1293            .await
1294            .unwrap();
1295        assert_eq!(command.action_intents.len(), 1);
1296        assert_eq!(command.action_intents[0].state, ActionState::Prepared);
1297        assert_eq!(command.action_intents[0].capability.id, "execute.demo");
1298        assert_eq!(command.action_intents[0].deadline, None);
1299
1300        let mut waiting = item("effect", command.next_state.clone());
1301        waiting.capability_pins = work.capability_pins.clone();
1302        let waiting_command = WorkflowDriver::<()>::evaluate(&driver, &(), &waiting)
1303            .await
1304            .unwrap();
1305        assert!(waiting_command.action_intents.is_empty());
1306        assert_eq!(waiting_command.event_type, "workflow.action_waiting");
1307
1308        waiting.wakeups.push(Wakeup {
1309            id: "terminal".into(),
1310            branch_id: None,
1311            kind: "timer".into(),
1312            payload: json!({
1313                "kind": "terminal_action",
1314                "observation": {
1315                    "id": "observation",
1316                    "action_intent_id": command.action_intents[0].id,
1317                    "provider_version": "1",
1318                    "observed_at": chrono::Utc::now(),
1319                    "state": "rejected",
1320                    "resource_ref": null,
1321                    "raw_receipt_digest": "receipt",
1322                    "terminal": true,
1323                    "retry_authorized": false
1324                }
1325            }),
1326        });
1327        let mut forged = waiting.clone();
1328        forged.wakeups[0].kind = "delivery".into();
1329        assert!(!forged.is_converging());
1330        let rejected = WorkflowDriver::<()>::evaluate(&driver, &(), &forged)
1331            .await
1332            .unwrap();
1333        assert_eq!(rejected.event_type, "workflow.action_waiting");
1334        assert!(rejected.consumed_wakeups.is_empty());
1335        waiting.wakeups.insert(
1336            0,
1337            Wakeup {
1338                id: "unrelated".into(),
1339                kind: "delivery".into(),
1340                branch_id: None,
1341                payload: json!({"value":999}),
1342            },
1343        );
1344        let terminal = WorkflowDriver::<()>::evaluate(&driver, &(), &waiting)
1345            .await
1346            .unwrap();
1347        assert_eq!(terminal.consumed_wakeups, [waiting.wakeups[1].id.clone()]);
1348        assert!(terminal.action_intents.is_empty());
1349        assert!(matches!(terminal.disposition, WorkDisposition::Complete));
1350        assert!(matches!(
1351            SpecDriver::new(&spec("ingress.cron", json!({})), &registry),
1352            Err(HostError::Schedule { .. })
1353        ));
1354    }
1355
1356    #[tokio::test]
1357    async fn delivery_does_not_advance_the_cron_cursor() {
1358        let registry = NodeRegistry::with_builtins();
1359        let driver = SpecDriver::new(
1360            &spec("ingress.cron", json!({"expression": "0 0 * * * *"})),
1361            &registry,
1362        )
1363        .unwrap();
1364        let mut work = item("counter", json!({}));
1365        work.scheduled_at = work.claimed_at + chrono::Duration::hours(1);
1366        work.wakeups.push(Wakeup {
1367            id: "delivery".into(),
1368            branch_id: None,
1369            kind: "delivery".into(),
1370            payload: json!({"value": 7}),
1371        });
1372        let command = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
1373            .await
1374            .unwrap();
1375        assert_eq!(command.next_state["state"]["__root__.seen"], json!([7]));
1376        assert!(matches!(
1377            command.disposition,
1378            WorkDisposition::Reschedule { at } if at == work.scheduled_at
1379        ));
1380    }
1381}