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