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