use std::collections::BTreeMap;
use std::sync::Arc;
use async_trait::async_trait;
use serde_json::{json, Map, Value};
use crate::{
ingress_plans, schedule_decision, ActionIntent, ActionState, CapabilityKind,
CapabilityManifest, CatchUpPolicy, EvaluationOutcome, Event, HostError, MemoryState,
NodeRegistry, ScheduleCadence, SchedulePolicy, Spec, Terminal, WorkDisposition, WorkItem,
WorkflowDriver, WorkflowHost, WorkflowTransitionCommand,
};
const IDLE_WAIT_SECS: i64 = 86_400;
pub struct SpecDriver {
spec_id: String,
host: WorkflowHost,
schedules: BTreeMap<String, SchedulePolicy>,
actions: BTreeMap<String, CapabilityManifest>,
}
#[derive(serde::Serialize, serde::Deserialize)]
struct ScheduleCursor {
next_at: chrono::DateTime<chrono::Utc>,
catch_up_count: u32,
}
impl SpecDriver {
pub fn new(spec: &Spec, registry: &NodeRegistry) -> Result<Self, HostError> {
let host = WorkflowHost::from_spec(spec, registry)?;
let schedules = ingress_plans(spec, registry)
.into_iter()
.map(|plan| {
let branch_id = plan.branch_id.clone();
(|| {
let required_u64 = |name: &str| {
plan.ingress_config[name]
.as_u64()
.ok_or_else(|| HostError::Schedule {
spec_id: spec.spec_id.clone(),
reason: format!("{} requires integer '{name}'", plan.ingress_type),
})
};
let cadence = match plan.ingress_type.as_str() {
"ingress.cron" => ScheduleCadence::Cron {
expression: plan.ingress_config["expression"]
.as_str()
.ok_or_else(|| HostError::Schedule {
spec_id: spec.spec_id.clone(),
reason: "ingress.cron requires a string 'expression'".into(),
})?
.to_owned(),
},
"ingress.fixed_rate" => ScheduleCadence::FixedRate {
milliseconds: required_u64("milliseconds")?,
},
"ingress.fixed_delay" => ScheduleCadence::FixedDelay {
milliseconds: required_u64("milliseconds")?,
},
_ => return Ok(None),
};
let catch_up = match plan.ingress_config["catch_up"].as_str().unwrap_or("once")
{
"skip" => CatchUpPolicy::Skip,
"once" => CatchUpPolicy::CatchUpOnce,
"all" => CatchUpPolicy::CatchUpAll {
limit: plan.ingress_config["catch_up_limit"]
.as_u64()
.and_then(|value| u32::try_from(value).ok())
.ok_or_else(|| HostError::Schedule {
spec_id: spec.spec_id.clone(),
reason: "catch_up=all requires positive integer catch_up_limit"
.into(),
})?,
},
value => {
return Err(HostError::Schedule {
spec_id: spec.spec_id.clone(),
reason: format!("unknown catch_up policy '{value}'"),
})
}
};
let policy = SchedulePolicy {
cadence,
timezone: plan.ingress_config["timezone"]
.as_str()
.unwrap_or("UTC")
.to_owned(),
catch_up,
};
policy.validate().map_err(|error| HostError::Schedule {
spec_id: spec.spec_id.clone(),
reason: error.to_string(),
})?;
Ok::<_, HostError>(Some(policy))
})()
.map(|policy| policy.map(|policy| (branch_id, policy)))
})
.collect::<Result<Vec<_>, _>>()?
.into_iter()
.flatten()
.collect();
let actions = registry
.capability_manifests()
.filter(|manifest| manifest.kind == CapabilityKind::Action)
.map(|manifest| (manifest.id.clone(), manifest.clone()))
.collect();
Ok(Self {
spec_id: spec.spec_id.clone(),
host,
schedules,
actions,
})
}
fn action_intents(
&self,
item: &WorkItem,
requests: Vec<crate::PreparedAction>,
) -> Result<Vec<ActionIntent>, String> {
requests
.into_iter()
.enumerate()
.map(|(index, request)| {
request.validate().map_err(|error| error.to_string())?;
let manifest = self
.actions
.get(request.capability_id.as_str())
.ok_or_else(|| {
format!(
"action step emitted unknown capability '{}'",
request.capability_id
)
})?;
let capability = item
.capability_pins
.iter()
.find(|pin| {
pin.id == manifest.id
&& pin.contract_version == manifest.contract_version
&& pin.content_digest == manifest.content_digest
})
.cloned()
.ok_or_else(|| {
format!(
"workflow revision does not pin action capability '{}'",
manifest.id
)
})?;
Ok(ActionIntent {
id: format!("{}:{}:action:{index}", item.id, item.state_version),
tenant_id: item.tenant_id.clone(),
instance_id: item
.id
.parse()
.map_err(|error: af_context::EmptyId| error.to_string())?,
run_id: item.run_id.clone(),
capability,
idempotency_key: format!(
"{}:{}:{}",
item.id, item.state_version, request.idempotency_key
),
state: ActionState::Prepared,
input: request.input,
effect: manifest.effect,
retry_class: manifest.idempotency_mode,
control_epochs: item.control_epochs,
resource_scope_id: request.resource_scope_id,
lease_epoch: item.lease_version,
action_epoch: item.state_version,
deadline: request.deadline,
reservation: request.reservation,
created_at: item.claimed_at,
})
})
.collect()
}
fn terminal_action(
&self,
item: &WorkItem,
next_state: &mut Map<String, Value>,
) -> Result<Option<WorkflowTransitionCommand>, String> {
let Some(wakeup) = item
.wakeups
.iter()
.find(|wakeup| wakeup.kind == "timer" && wakeup.payload["kind"] == "terminal_action")
else {
return Ok(None);
};
let observation: crate::ActionObservation =
serde_json::from_value(wakeup.payload["observation"].clone())
.map_err(|error| format!("terminal action fact: {error}"))?;
let Some(internal) = next_state
.get_mut("__workflow")
.and_then(Value::as_object_mut)
else {
return Err("terminal action has no durable pending-action state".into());
};
let pending = internal
.get_mut("pending_actions")
.and_then(Value::as_array_mut)
.ok_or_else(|| "terminal action has no pending action list".to_string())?;
let before = pending.len();
pending.retain(|id| id.as_str() != Some(&observation.action_intent_id));
if pending.len() == before {
return Err("terminal action does not belong to this workflow instance".into());
}
let disposition = if pending.is_empty() {
internal
.remove("resume_at")
.map(serde_json::from_value)
.transpose()
.map_err(|error| format!("stored workflow resume_at: {error}"))?
.map(|at| WorkDisposition::Reschedule { at })
.unwrap_or_else(|| {
if internal
.get("static_event_consumed")
.and_then(Value::as_bool)
.unwrap_or(false)
{
WorkDisposition::Complete
} else {
Self::idle_disposition()
}
})
} else {
WorkDisposition::Continue {
delay_secs: IDLE_WAIT_SECS,
}
};
let succeeded = matches!(
observation.state.as_str(),
"completed" | "max_steps_reached" | "succeeded"
);
let mut selected = item.clone();
selected.wakeups = vec![wakeup.clone()];
Ok(Some(command(
&selected,
Value::Object(next_state.clone()),
Vec::new(),
EvaluationOutcome {
triggered: true,
matched: true,
succeeded,
action_terminal: true,
},
disposition,
"workflow.action_observed",
)))
}
fn catch_up_count(state: &Map<String, Value>) -> u32 {
state
.get("__workflow")
.and_then(Value::as_object)
.and_then(|internal| internal.get("catch_up_count"))
.and_then(Value::as_u64)
.and_then(|value| u32::try_from(value).ok())
.unwrap_or(0)
}
fn set_internal(state: &mut Map<String, Value>, key: &str, value: Value) -> Result<(), String> {
let internal = state
.entry("__workflow")
.or_insert_with(|| json!({}))
.as_object_mut()
.ok_or_else(|| "workflow internal state must be an object".to_string())?;
internal.insert(key.into(), value);
Ok(())
}
fn schedule_cursors(
&self,
item: &WorkItem,
state: &Map<String, Value>,
) -> Result<BTreeMap<String, ScheduleCursor>, String> {
let saved = state
.get("__workflow")
.and_then(|value| value.get("schedules"));
let mut cursors: BTreeMap<String, ScheduleCursor> = saved
.cloned()
.map(serde_json::from_value)
.transpose()
.map_err(|error| format!("stored branch schedule: {error}"))?
.unwrap_or_default();
for (branch, policy) in &self.schedules {
if !cursors.contains_key(branch) {
if saved.is_some() {
return Err("stored schedules omit a pinned branch".into());
}
let legacy_root = self.host.branch_count() == 1
|| (saved.is_none()
&& item.state_version > 0
&& branch == self.host.branch_id()
&& (state.contains_key("last_run")
|| state.get("__workflow").is_some_and(|internal| {
internal.get("catch_up_count").is_some()
|| internal.get("resume_at").is_some()
})));
let next_at = if legacy_root {
state
.get("__workflow")
.and_then(|internal| internal.get("resume_at"))
.cloned()
.map(serde_json::from_value)
.transpose()
.map_err(|error| format!("stored workflow resume_at: {error}"))?
.unwrap_or(item.scheduled_at)
} else {
crate::next_scheduled_at(
policy,
item.lifecycle
.starts_at
.map_or(item.created_at, |starts| starts.max(item.created_at)),
None,
)?
};
cursors.insert(
branch.clone(),
ScheduleCursor {
next_at,
catch_up_count: if legacy_root {
Self::catch_up_count(state)
} else {
0
},
},
);
}
}
if cursors
.keys()
.any(|branch| !self.schedules.contains_key(branch))
{
return Err("stored schedule refers to an unknown branch".into());
}
Ok(cursors)
}
fn idle_disposition() -> WorkDisposition {
WorkDisposition::Continue {
delay_secs: IDLE_WAIT_SECS,
}
}
}
#[async_trait]
impl<Context: Send + Sync> WorkflowDriver<Context> for SpecDriver {
fn name(&self) -> &'static str {
"spec"
}
fn spec_ids(&self) -> Vec<&str> {
vec![&self.spec_id]
}
fn validate_specs(&self) -> Result<(), String> {
Ok(())
}
async fn evaluate(
&self,
_: &Context,
item: &WorkItem,
) -> Result<WorkflowTransitionCommand, String> {
if item.spec_id != self.spec_id {
return Err("compiled spec does not match the claimed spec identity".into());
}
let mut next_state = match &item.config {
Value::Object(map) => map.clone(),
Value::Null => Map::new(),
_ => return Err("spec instance config must be a JSON object".into()),
};
if item.cancel_requested {
return Ok(command(
item,
Value::Object(next_state),
Vec::new(),
EvaluationOutcome::default(),
WorkDisposition::Complete,
"workflow.cancelled",
));
}
if let Some(command) = self.terminal_action(item, &mut next_state)? {
return Ok(command);
}
if next_state
.get("__workflow")
.and_then(Value::as_object)
.and_then(|internal| internal.get("pending_actions"))
.and_then(Value::as_array)
.is_some_and(|pending| !pending.is_empty())
{
return Ok(command(
item,
Value::Object(next_state),
Vec::new(),
EvaluationOutcome::default(),
Self::idle_disposition(),
"workflow.action_waiting",
));
}
let mut cursors = self.schedule_cursors(item, &next_state)?;
let due = cursors
.values()
.any(|cursor| cursor.next_at <= item.claimed_at);
let prefer_schedule = due
&& !item.wakeups.is_empty()
&& next_state
.get("__workflow")
.and_then(|internal| internal.get("last_dispatch"))
.and_then(Value::as_str)
!= Some("schedule");
let mut selected = item.clone();
if prefer_schedule {
selected.wakeups.clear();
}
let item = &selected;
let earliest = cursors
.iter()
.min_by_key(|(branch, cursor)| (cursor.next_at, *branch));
let branch_id = if let Some(wakeup) = item.wakeups.first() {
wakeup
.branch_id
.as_ref()
.map(|id| id.as_str())
.or_else(|| (self.host.branch_count() == 1).then(|| self.host.branch_id()))
.ok_or_else(|| "multi-branch wakeup requires branch_id".to_string())?
.to_owned()
} else if let Some((branch, _)) = earliest {
branch.clone()
} else if self.host.branch_count() == 1 {
self.host.branch_id().to_owned()
} else {
return Err("multi-branch workflow requires a branch-addressed delivery".into());
};
if !self.host.branch_ids().any(|branch| branch == branch_id) {
return Err("wakeup refers to an unknown branch".into());
}
let schedule = if item.wakeups.is_empty() {
if let Some(cursor) = cursors.get_mut(&branch_id) {
let decision = if cursor.next_at > item.claimed_at {
crate::ScheduleDecision {
tick_at: None,
next_at: cursor.next_at,
catch_up_count: cursor.catch_up_count,
}
} else {
schedule_decision(
&self.schedules[&branch_id],
cursor.next_at,
item.claimed_at,
cursor.catch_up_count,
)?
};
cursor.next_at = decision.next_at;
cursor.catch_up_count = decision.catch_up_count;
Some(decision)
} else {
None
}
} else {
None
};
Self::set_internal(
&mut next_state,
"last_dispatch",
json!(if item.wakeups.is_empty() {
"schedule"
} else {
"delivery"
}),
)?;
Self::set_internal(&mut next_state, "schedules", json!(cursors))?;
let next_at = cursors.values().map(|cursor| cursor.next_at).min();
if let Some(decision) = schedule {
if decision.tick_at.is_none() {
return Ok(command(
item,
Value::Object(next_state),
Vec::new(),
EvaluationOutcome::default(),
WorkDisposition::Reschedule {
at: next_at.expect("selected schedule has a cursor"),
},
"workflow.schedule_skipped",
));
}
}
let state = Arc::new(MemoryState::from_snapshot(
next_state
.get("state")
.and_then(Value::as_object)
.cloned()
.unwrap_or_default(),
));
let wakeup_payload = item.wakeups.first().map(|wakeup| wakeup.payload.clone());
let static_event = next_state.get("event").cloned().filter(|_| {
!next_state
.get("__workflow")
.and_then(Value::as_object)
.and_then(|internal| internal.get("static_event_consumed"))
.and_then(Value::as_bool)
.unwrap_or(false)
});
let consumed_static_event = wakeup_payload.is_none() && static_event.is_some();
let payload = wakeup_payload
.or(static_event)
.or_else(|| {
schedule
.and_then(|decision| decision.tick_at)
.map(|tick_at| json!({ "tick_at": tick_at }))
})
.ok_or_else(|| "event workflow requires a trigger delivery".to_string())?;
if consumed_static_event {
Self::set_internal(&mut next_state, "static_event_consumed", json!(true))?;
}
let context = self.host.context_for(&branch_id, state.clone());
let outcome = self
.host
.run_event(&context, Event::from_json(payload))
.await
.ok_or_else(|| "spec has no selected branch".to_string())?;
let (terminal, exit_reason) = match &outcome.terminal {
Terminal::Completed => ("completed", Value::Null),
Terminal::Dropped { node_id, reason } => {
("dropped", json!({ "node_id": node_id, "reason": reason }))
}
};
next_state.insert("state".into(), Value::Object(state.snapshot()));
next_state.insert(
"last_run".into(),
json!({
"branch_id": branch_id,
"terminal": terminal,
"exit": exit_reason,
"steps_run": outcome.steps_run,
"survivors": outcome.survivors.len(),
"at": item.claimed_at,
}),
);
let action_intents = self.action_intents(item, outcome.actions)?;
let disposition = match next_at {
Some(at) if action_intents.is_empty() => WorkDisposition::Reschedule { at },
Some(at) => {
Self::set_internal(&mut next_state, "resume_at", json!(at))?;
Self::idle_disposition()
}
None if action_intents.is_empty() && consumed_static_event => WorkDisposition::Complete,
None => Self::idle_disposition(),
};
if !action_intents.is_empty() {
let internal = next_state
.entry("__workflow")
.or_insert_with(|| json!({}))
.as_object_mut()
.ok_or_else(|| "workflow internal state must be an object".to_string())?;
let pending = internal
.entry("pending_actions")
.or_insert_with(|| json!([]))
.as_array_mut()
.ok_or_else(|| "workflow pending actions must be an array".to_string())?;
pending.extend(action_intents.iter().map(|intent| json!(intent.id)));
}
let evaluation = EvaluationOutcome {
triggered: true,
matched: outcome.matched,
succeeded: outcome.succeeded && action_intents.is_empty(),
action_terminal: false,
};
Ok(command(
item,
Value::Object(next_state),
action_intents,
evaluation,
disposition,
"workflow.spec_evaluated",
))
}
}
fn command(
item: &WorkItem,
next_state: Value,
action_intents: Vec<ActionIntent>,
outcome: EvaluationOutcome,
disposition: WorkDisposition,
event_type: &str,
) -> WorkflowTransitionCommand {
WorkflowTransitionCommand {
consumed_wakeups: if outcome.triggered {
item.wakeups
.first()
.map(|wakeup| wakeup.id.clone())
.into_iter()
.collect()
} else {
Vec::new()
},
delivery_key: format!("spec:{}:{}", item.id, item.state_version),
delivery_digest: format!(
"{}:{}:{}",
item.workflow_revision_digest, item.execution_profile_digest, item.state_version
),
event_type: event_type.into(),
event_digest: format!("{event_type}:{}:{}", item.id, item.state_version),
event_payload: json!({ "spec_id": item.spec_id, "wakeups": item.wakeups.len() }),
next_state,
action_intents,
outcome,
disposition,
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::{ControlEpochs, DriverRegistry, PreparedAction, StepNode, StepResult, Wakeup};
struct PassNode;
#[async_trait::async_trait]
impl StepNode for PassNode {
async fn process(&self, event: &Event, _: &crate::WorkflowContext) -> StepResult {
StepResult::Pass(event.clone())
}
}
struct ActionNode;
#[async_trait::async_trait]
impl StepNode for ActionNode {
async fn process(&self, event: &Event, _: &crate::WorkflowContext) -> StepResult {
StepResult::Action {
event: event.clone(),
action: Box::new(PreparedAction::new(
"execute.demo".parse().unwrap(),
"request-1",
json!({"value": 1}),
)),
}
}
}
struct VersionNode(&'static str);
#[async_trait::async_trait]
impl StepNode for VersionNode {
async fn process(&self, event: &Event, ctx: &crate::WorkflowContext) -> StepResult {
ctx.state
.append(&ctx.scoped("seen"), json!(self.0), Some(3), None)
.await;
StepResult::Pass(event.clone())
}
}
fn version_one(_: &Value) -> Result<Box<dyn StepNode>, crate::NodeError> {
Ok(Box::new(VersionNode("v1")))
}
fn version_two(_: &Value) -> Result<Box<dyn StepNode>, crate::NodeError> {
Ok(Box::new(VersionNode("v2")))
}
fn item(spec_id: &str, config: Value) -> WorkItem {
WorkItem {
id: "instance".into(),
run_id: uuid::Uuid::new_v4().to_string().parse().unwrap(),
tenant_id: "tenant".parse().unwrap(),
subject_id: "subject".parse().unwrap(),
spec_id: spec_id.into(),
definition_id: spec_id.into(),
workflow_revision: 1,
workflow_revision_digest: "digest".into(),
execution_profile_id: "profile".into(),
execution_profile_revision: 1,
execution_profile_digest: "profile-digest".into(),
kernel_abi_version: "1".into(),
capability_pins: Vec::new(),
lifecycle: crate::LifecyclePolicy::run_once(),
scheduled_at: chrono::Utc::now(),
created_at: chrono::Utc::now(),
claimed_at: chrono::Utc::now(),
config,
state_version: 0,
control_epochs: ControlEpochs::default(),
cancel_requested: false,
lease_version: 1,
wakeups: Vec::new(),
}
}
fn spec(ingress: &str, ingress_config: Value) -> Spec {
Spec::from_json(
&json!({
"spec_id": "counter", "version": "1",
"branches": [{
"branch_id": "__root__",
"nodes": [
{"id": "in", "type": ingress, "config": ingress_config},
{"id": "count", "type": "transform.state_append",
"config": {"key": "seen", "path": "value", "max_len": 3}}
],
"edges": [{"source": "in", "target": "count"}]
}]
})
.to_string(),
)
.unwrap()
}
#[tokio::test]
async fn published_revisions_select_distinct_graphs_and_reject_conflicting_pins() {
let nodes = NodeRegistry::with_builtins();
let mut registry = crate::DriverRegistry::<()>::new().with_node_registry(nodes);
let first_spec = spec("ingress.event", json!({}));
registry.prewarm_spec(&first_spec).unwrap();
let first = item("counter", json!({"event": {"value": 7}}));
let revision = crate::WorkflowRevision {
definition_id: first.definition_id.clone(),
revision: 1,
content_digest: first.workflow_revision_digest.clone(),
kernel_abi_version: "1".into(),
dependency_set_digest: "deps".into(),
expression_versions: Default::default(),
capabilities: Vec::new(),
template_provenance: Value::Null,
spec: first_spec,
};
let first_driver = registry
.resolve_revision(&first, revision.clone())
.await
.unwrap();
let cached = registry
.resolve_revision(&first, revision.clone())
.await
.unwrap();
assert!(Arc::ptr_eq(&first_driver, &cached));
let mut second = first.clone();
second.workflow_revision = 2;
second.workflow_revision_digest = "revision-2".into();
let mut second_revision = revision.clone();
second_revision.revision = 2;
second_revision.content_digest = second.workflow_revision_digest.clone();
second_revision.spec.branches[0].nodes[1].config["key"] = json!("second");
let second_driver = registry
.resolve_revision(&second, second_revision.clone())
.await
.unwrap();
let first_result = first_driver.evaluate(&(), &first).await.unwrap();
let second_result = second_driver.evaluate(&(), &second).await.unwrap();
assert_eq!(
first_result.next_state["state"]["__root__.seen"],
json!([7])
);
assert_eq!(
second_result.next_state["state"]["__root__.second"],
json!([7])
);
assert!(second_result.next_state["state"]
.get("__root__.seen")
.is_none());
assert!(registry
.resolve_revision(&first, second_revision.clone())
.await
.is_err());
second_revision.spec.branches[0].nodes[1].config["key"] = json!("conflict");
assert!(registry
.resolve_revision(&second, second_revision.clone())
.await
.is_err());
let mut other_tenant = second.clone();
other_tenant.tenant_id = "other-tenant".parse().unwrap();
assert!(registry
.resolve_revision(&other_tenant, second_revision.clone())
.await
.is_ok());
let mut other_definition = second.clone();
other_definition.definition_id = "other-definition".into();
second_revision.definition_id = other_definition.definition_id.clone();
assert!(registry
.resolve_revision(&other_definition, second_revision)
.await
.is_ok());
let mut unsupported = revision.clone();
unsupported
.expression_versions
.insert("transform.state_append".into(), "unknown".into());
assert!(registry
.resolve_revision(&first, unsupported)
.await
.is_err());
let mut guarded_nodes = NodeRegistry::with_builtins();
let mut manifest = CapabilityManifest::action(
"transform.state_append",
"1",
"read-v1",
crate::Effect::Read,
crate::IdempotencyMode::None,
false,
);
manifest.kind = CapabilityKind::Expression;
guarded_nodes.register_capability(manifest).unwrap();
let guarded = crate::DriverRegistry::<()>::new().with_node_registry(guarded_nodes);
assert!(
guarded
.resolve_revision(&first, revision.clone())
.await
.is_err(),
"read capabilities need pins even though they do not emit actions"
);
for version in 3..=259 {
let mut next_item = first.clone();
next_item.workflow_revision = version;
next_item.workflow_revision_digest = format!("revision-{version}");
let mut next_revision = revision.clone();
next_revision.revision = version;
next_revision.content_digest = next_item.workflow_revision_digest.clone();
registry
.resolve_revision(&next_item, next_revision)
.await
.unwrap();
}
let (left, right) = tokio::join!(
registry.resolve_revision(&first, revision.clone()),
registry.resolve_revision(&first, revision.clone())
);
let (left, right) = (left.unwrap(), right.unwrap());
assert!(
!Arc::ptr_eq(&left, &first_driver),
"oldest compiled revision is evicted"
);
assert!(
Arc::ptr_eq(&left, &right),
"concurrent resolution shares one compilation"
);
let replaced = registry.with_node_registry(NodeRegistry::empty());
assert!(
replaced
.resolve_revision(&first, revision.clone())
.await
.is_err(),
"changing the node registry must invalidate compiled graphs"
);
let mut unavailable = revision;
unavailable.capabilities.push(crate::CapabilityPin {
id: "missing".into(),
contract_version: "1".into(),
content_digest: "missing-v1".into(),
});
let mut unavailable_item = first;
unavailable_item.capability_pins = unavailable.capabilities.clone();
assert!(replaced
.resolve_revision(&unavailable_item, unavailable)
.await
.is_err());
}
#[tokio::test]
async fn published_revisions_keep_two_versions_of_one_capability_runnable() {
let capability = |version: &str, digest: &str| {
let mut manifest = CapabilityManifest::action(
"transform.state_append",
version,
digest,
crate::Effect::Read,
crate::IdempotencyMode::None,
false,
);
manifest.kind = CapabilityKind::Expression;
manifest
};
let v1 = capability("1", "state-append-v1");
let v2 = capability("2", "state-append-v2");
let pin = |manifest: &CapabilityManifest| crate::CapabilityPin {
id: manifest.id.clone(),
contract_version: manifest.contract_version.clone(),
content_digest: manifest.content_digest.clone(),
};
let mut nodes = NodeRegistry::with_builtins();
let schema = nodes.schema("transform.state_append").cloned().unwrap();
nodes.register_capability(v1.clone()).unwrap();
nodes.register_capability(v2.clone()).unwrap();
nodes
.register_capability_implementation(
pin(&v1),
version_one,
Some(schema.clone()),
None,
false,
)
.unwrap();
nodes
.register_capability_implementation(pin(&v2), version_two, Some(schema), None, false)
.unwrap();
let registry = crate::DriverRegistry::<()>::new().with_node_registry(nodes);
let mut first = item("counter", json!({"event": {"value": 1}}));
first.capability_pins = vec![pin(&v1)];
let mut first_revision = crate::WorkflowRevision {
definition_id: first.definition_id.clone(),
revision: first.workflow_revision,
content_digest: first.workflow_revision_digest.clone(),
kernel_abi_version: "1".into(),
dependency_set_digest: "deps-v1".into(),
expression_versions: Default::default(),
capabilities: first.capability_pins.clone(),
template_provenance: Value::Null,
spec: spec("ingress.event", json!({})),
};
let mut second = first.clone();
second.workflow_revision = 2;
second.workflow_revision_digest = "revision-2".into();
second.capability_pins = vec![pin(&v2)];
let mut second_revision = first_revision.clone();
second_revision.revision = 2;
second_revision.content_digest = second.workflow_revision_digest.clone();
second_revision.dependency_set_digest = "deps-v2".into();
second_revision.capabilities = second.capability_pins.clone();
let first_driver = registry
.resolve_revision(&first, first_revision.clone())
.await
.unwrap();
let second_driver = registry
.resolve_revision(&second, second_revision)
.await
.unwrap();
assert_eq!(
first_driver.evaluate(&(), &first).await.unwrap().next_state["state"]["__root__.seen"],
json!(["v1"])
);
assert_eq!(
second_driver
.evaluate(&(), &second)
.await
.unwrap()
.next_state["state"]["__root__.seen"],
json!(["v2"])
);
first_revision.capabilities = vec![pin(&v2)];
assert!(registry
.resolve_revision(&first, first_revision)
.await
.is_err());
}
#[tokio::test]
async fn old_root_cursor_is_preserved_and_due_schedules_are_not_starved() {
let mut graph = spec("ingress.fixed_rate", json!({"milliseconds": 10_000}));
let mut sibling = spec("ingress.event", json!({})).branches.remove(0);
sibling.branch_id = "events".into();
for node in &mut sibling.nodes {
node.id = format!("events-{}", node.id);
}
for edge in &mut sibling.edges {
edge.source = format!("events-{}", edge.source);
edge.target = format!("events-{}", edge.target);
}
graph.branches.push(sibling);
let driver = SpecDriver::new(&graph, &NodeRegistry::with_builtins()).unwrap();
let start: chrono::DateTime<chrono::Utc> = "2026-01-01T00:00:00Z".parse().unwrap();
let mut work = item("counter", json!({"__workflow":{"catch_up_count":2}}));
work.created_at = start;
work.state_version = 7;
work.scheduled_at = start + chrono::Duration::seconds(100);
work.claimed_at = work.scheduled_at;
let cursors = driver
.schedule_cursors(&work, work.config.as_object().unwrap())
.unwrap();
assert_eq!(cursors["__root__"].next_at, work.scheduled_at);
assert_eq!(cursors["__root__"].catch_up_count, 2);
work.config["__workflow"]["resume_at"] = json!(start + chrono::Duration::seconds(150));
let cursors = driver
.schedule_cursors(&work, work.config.as_object().unwrap())
.unwrap();
assert_eq!(
cursors["__root__"].next_at,
start + chrono::Duration::seconds(150)
);
work.config["__workflow"]
.as_object_mut()
.unwrap()
.remove("resume_at");
work.wakeups.push(serde_json::from_value(json!({"id":"queued", "kind":"delivery", "branch_id":"events", "payload":{"value":7}})).unwrap());
let schedule = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
.await
.unwrap();
assert_eq!(schedule.next_state["last_run"]["branch_id"], "__root__");
assert!(
schedule.consumed_wakeups.is_empty(),
"schedule does not acknowledge queued event"
);
work.config = schedule.next_state;
let event = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
.await
.unwrap();
assert_eq!(event.next_state["last_run"]["branch_id"], "events");
assert_eq!(event.consumed_wakeups, ["queued"]);
work.config = event.next_state;
work.claimed_at += chrono::Duration::seconds(10);
let schedule = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
.await
.unwrap();
assert_eq!(schedule.next_state["last_run"]["branch_id"], "__root__");
assert!(schedule.consumed_wakeups.is_empty());
}
#[test]
fn first_failure_does_not_migrate_a_nonexistent_root_cursor() {
let mut graph = spec("ingress.fixed_rate", json!({"milliseconds":100_000}));
let mut fast = graph.branches[0].clone();
fast.branch_id = "fast".into();
fast.nodes[0].config["milliseconds"] = json!(10_000);
for node in &mut fast.nodes {
node.id = format!("fast-{}", node.id);
}
for edge in &mut fast.edges {
edge.source = format!("fast-{}", edge.source);
edge.target = format!("fast-{}", edge.target);
}
graph.branches.push(fast);
let driver = SpecDriver::new(&graph, &NodeRegistry::with_builtins()).unwrap();
let mut work = item("counter", json!({}));
work.created_at = "2026-01-01T00:00:00Z".parse().unwrap();
work.scheduled_at = work.created_at + chrono::Duration::seconds(10);
work.state_version = 1;
let cursors = driver
.schedule_cursors(&work, work.config.as_object().unwrap())
.unwrap();
assert_eq!(
cursors["__root__"].next_at,
work.created_at + chrono::Duration::seconds(100)
);
assert_eq!(cursors["fast"].next_at, work.scheduled_at);
}
#[tokio::test]
async fn branches_keep_independent_schedule_cursors_across_restart() {
let mut graph = spec("ingress.event", json!({}));
for (branch_id, period) in [("fast", 10_000), ("slow", 25_000)] {
let mut branch = spec("ingress.fixed_rate", json!({"milliseconds": period}))
.branches
.remove(0);
branch.branch_id = branch_id.into();
for node in &mut branch.nodes {
node.id = format!("{branch_id}-{}", node.id);
}
for edge in &mut branch.edges {
edge.source = format!("{branch_id}-{}", edge.source);
edge.target = format!("{branch_id}-{}", edge.target);
}
branch.nodes[1].config["path"] = json!("tick_at");
graph.branches.push(branch);
}
let start: chrono::DateTime<chrono::Utc> = "2026-01-01T00:00:00Z".parse().unwrap();
let mut work = item("counter", json!({}));
work.created_at = start;
work.scheduled_at = start + chrono::Duration::seconds(10);
for (second, branch) in [(10, "fast"), (20, "fast"), (25, "slow"), (30, "fast")] {
work.claimed_at = start + chrono::Duration::seconds(second);
let restarted = SpecDriver::new(&graph, &NodeRegistry::with_builtins()).unwrap();
let command = WorkflowDriver::<()>::evaluate(&restarted, &(), &work)
.await
.unwrap();
assert_eq!(command.next_state["last_run"]["branch_id"], branch);
let WorkDisposition::Reschedule { at } = command.disposition else {
panic!("scheduled branch must retain every cursor");
};
work.config = serde_json::from_str(&command.next_state.to_string()).unwrap();
work.scheduled_at = at;
}
assert_eq!(
work.config["state"]["fast.seen"].as_array().unwrap().len(),
3
);
assert_eq!(
work.config["state"]["slow.seen"],
json!([start + chrono::Duration::seconds(25)])
);
let saved_cursors = work.config["__workflow"]["schedules"].clone();
work.wakeups.push(serde_json::from_value(json!({"id":"event", "kind":"delivery", "branch_id":"__root__", "payload":{"value":99}})).unwrap());
let driver = SpecDriver::new(&graph, &NodeRegistry::with_builtins()).unwrap();
let command = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
.await
.unwrap();
assert_eq!(command.next_state["__workflow"]["schedules"], saved_cursors);
assert_eq!(command.next_state["state"]["__root__.seen"], json!([99]));
work.wakeups.clear();
work.claimed_at = start + chrono::Duration::seconds(31);
let skipped = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
.await
.unwrap();
assert_eq!(skipped.event_type, "workflow.schedule_skipped");
assert_eq!(skipped.next_state["state"], work.config["state"]);
}
#[tokio::test]
async fn delivery_selects_branch_and_rejects_ambiguous_legacy_identity() {
let mut graph = spec("ingress.event", json!({}));
let mut sibling = graph.branches[0].clone();
sibling.branch_id = "other".into();
for node in &mut sibling.nodes {
node.id = format!("other-{}", node.id);
}
for edge in &mut sibling.edges {
edge.source = format!("other-{}", edge.source);
edge.target = format!("other-{}", edge.target);
}
graph.branches.push(sibling);
let driver = SpecDriver::new(&graph, &NodeRegistry::with_builtins()).unwrap();
let mut work = item("counter", json!({}));
work.wakeups.push(
serde_json::from_value(json!({
"id": "delivery", "kind": "delivery", "branch_id": "other",
"payload": {"value": 7}
}))
.unwrap(),
);
let command = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
.await
.unwrap();
assert_eq!(command.next_state["state"]["other.seen"], json!([7]));
assert!(command.next_state["state"]["__root__.seen"].is_null());
work.config = command.next_state;
work.wakeups[0] = serde_json::from_value(json!({
"id": "next", "kind": "delivery", "branch_id": "__root__",
"payload": {"value": 9}
}))
.unwrap();
let restarted = SpecDriver::new(&graph, &NodeRegistry::with_builtins()).unwrap();
let command = WorkflowDriver::<()>::evaluate(&restarted, &(), &work)
.await
.unwrap();
assert_eq!(command.next_state["state"]["other.seen"], json!([7]));
assert_eq!(command.next_state["state"]["__root__.seen"], json!([9]));
for branch in [Value::Null, json!("missing")] {
work.wakeups[0] = serde_json::from_value(json!({
"id": "bad", "kind": "delivery", "branch_id": branch,
"payload": {"value": 1}
}))
.unwrap();
assert!(WorkflowDriver::<()>::evaluate(&driver, &(), &work)
.await
.is_err());
}
}
#[tokio::test]
async fn compiled_driver_rejects_another_spec_identity() {
let registry = NodeRegistry::with_builtins();
let driver = SpecDriver::new(&spec("ingress.event", json!({})), ®istry).unwrap();
let result = WorkflowDriver::<()>::evaluate(
&driver,
&(),
&item("different-spec", json!({"event": {"value": 1}})),
)
.await;
assert!(
result.is_err(),
"a compiled graph must not execute another spec"
);
}
#[tokio::test]
async fn cron_spec_reschedules_and_persists_branch_state() {
let registry = NodeRegistry::with_builtins();
let driver = SpecDriver::new(
&spec("ingress.cron", json!({"expression": "0 0 * * * *"})),
®istry,
)
.unwrap();
let first = WorkflowDriver::<()>::evaluate(
&driver,
&(),
&item("counter", json!({"event": {"value": 1}})),
)
.await
.unwrap();
let WorkDisposition::Reschedule { at } = first.disposition else {
panic!("cron spec must reschedule");
};
assert_eq!(first.next_state["last_run"]["terminal"], "completed");
let mut next = item("counter", first.next_state);
next.scheduled_at = at;
next.claimed_at = at;
let second = WorkflowDriver::<()>::evaluate(&driver, &(), &next)
.await
.unwrap();
assert_eq!(
second.next_state["state"]["__root__.seen"],
json!([1]),
"the consumed bootstrap event must not be replayed on a cron tick"
);
assert_eq!(second.next_state["last_run"]["at"], json!(at));
}
#[tokio::test]
async fn event_spec_consumes_a_wakeup_then_waits() {
let registry = NodeRegistry::with_builtins();
let driver = SpecDriver::new(&spec("ingress.event", json!({})), ®istry).unwrap();
assert_eq!(
WorkflowDriver::<()>::evaluate(&driver, &(), &item("counter", json!({})))
.await
.unwrap_err(),
"event workflow requires a trigger delivery"
);
let mut work = item("counter", json!({}));
work.wakeups.push(Wakeup {
id: "delivery".into(),
branch_id: None,
kind: "delivery".into(),
payload: json!({"value": 7}),
});
let mut second_wakeup = work.wakeups[0].clone();
second_wakeup.id = "later-delivery".into();
second_wakeup.payload = json!({"value":8});
work.wakeups.push(second_wakeup);
let command = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
.await
.unwrap();
assert_eq!(command.consumed_wakeups, ["delivery"]);
assert_eq!(command.next_state["state"]["__root__.seen"], json!([7]));
assert!(matches!(
command.disposition,
WorkDisposition::Continue { .. }
));
let mut registry = DriverRegistry::<()>::new();
registry.register(Arc::new(driver)).unwrap();
assert_eq!(registry.spec_ids(), ["counter"]);
}
#[tokio::test]
async fn effectful_steps_are_lifted_into_pinned_intents() {
let mut registry = NodeRegistry::with_builtins();
registry.register_step("guard.auth", |_| Ok(Box::new(PassNode)));
registry.register_side_effect_guard("guard.auth");
registry.register_step("execute.demo", |_| Ok(Box::new(ActionNode)));
let manifest = CapabilityManifest::action(
"execute.demo",
"1",
"demo-digest",
crate::Effect::ExternalWrite,
crate::IdempotencyMode::Native,
true,
);
registry.register_capability(manifest.clone()).unwrap();
let effect = Spec::from_json(
&json!({
"spec_id": "effect", "version": "1",
"branches": [{
"branch_id": "__root__",
"nodes": [
{"id": "in", "type": "ingress.event", "config": {}},
{"id": "auth", "type": "guard.auth", "config": {}},
{"id": "do", "type": "execute.demo", "config": {}}
],
"edges": [
{"source": "in", "target": "auth"},
{"source": "auth", "target": "do"}
]
}]
})
.to_string(),
)
.unwrap();
let driver = SpecDriver::new(&effect, ®istry).unwrap();
let mut work = item("effect", json!({"event": {"value": 1}}));
work.capability_pins.push(crate::CapabilityPin {
id: manifest.id,
contract_version: manifest.contract_version,
content_digest: manifest.content_digest,
});
let command = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
.await
.unwrap();
assert_eq!(command.action_intents.len(), 1);
assert_eq!(command.action_intents[0].state, ActionState::Prepared);
assert_eq!(command.action_intents[0].capability.id, "execute.demo");
assert_eq!(command.action_intents[0].deadline, None);
let mut waiting = item("effect", command.next_state.clone());
waiting.capability_pins = work.capability_pins.clone();
let waiting_command = WorkflowDriver::<()>::evaluate(&driver, &(), &waiting)
.await
.unwrap();
assert!(waiting_command.action_intents.is_empty());
assert_eq!(waiting_command.event_type, "workflow.action_waiting");
waiting.wakeups.push(Wakeup {
id: "terminal".into(),
branch_id: None,
kind: "timer".into(),
payload: json!({
"kind": "terminal_action",
"observation": {
"id": "observation",
"action_intent_id": command.action_intents[0].id,
"provider_version": "1",
"observed_at": chrono::Utc::now(),
"state": "rejected",
"resource_ref": null,
"raw_receipt_digest": "receipt",
"terminal": true,
"retry_authorized": false
}
}),
});
let mut forged = waiting.clone();
forged.wakeups[0].kind = "delivery".into();
assert!(!forged.is_converging());
let rejected = WorkflowDriver::<()>::evaluate(&driver, &(), &forged)
.await
.unwrap();
assert_eq!(rejected.event_type, "workflow.action_waiting");
assert!(rejected.consumed_wakeups.is_empty());
waiting.wakeups.insert(
0,
Wakeup {
id: "unrelated".into(),
kind: "delivery".into(),
branch_id: None,
payload: json!({"value":999}),
},
);
let terminal = WorkflowDriver::<()>::evaluate(&driver, &(), &waiting)
.await
.unwrap();
assert_eq!(terminal.consumed_wakeups, [waiting.wakeups[1].id.clone()]);
assert!(terminal.action_intents.is_empty());
assert!(matches!(terminal.disposition, WorkDisposition::Complete));
assert!(matches!(
SpecDriver::new(&spec("ingress.cron", json!({})), ®istry),
Err(HostError::Schedule { .. })
));
}
#[tokio::test]
async fn delivery_does_not_advance_the_cron_cursor() {
let registry = NodeRegistry::with_builtins();
let driver = SpecDriver::new(
&spec("ingress.cron", json!({"expression": "0 0 * * * *"})),
®istry,
)
.unwrap();
let mut work = item("counter", json!({}));
work.scheduled_at = work.claimed_at + chrono::Duration::hours(1);
work.wakeups.push(Wakeup {
id: "delivery".into(),
branch_id: None,
kind: "delivery".into(),
payload: json!({"value": 7}),
});
let command = WorkflowDriver::<()>::evaluate(&driver, &(), &work)
.await
.unwrap();
assert_eq!(command.next_state["state"]["__root__.seen"], json!([7]));
assert!(matches!(
command.disposition,
WorkDisposition::Reschedule { at } if at == work.scheduled_at
));
}
}