1use 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
20const IDLE_WAIT_SECS: i64 = 86_400;
24
25pub 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 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 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 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 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!({})), ®istry).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 ®istry,
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!({})), ®istry).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, ®istry).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!({})), ®istry),
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 ®istry,
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}