use std::collections::HashMap;
use std::sync::{Arc, Mutex};
use std::time::Duration;
use chrono::{TimeDelta, Utc};
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
use serde_json::{Value, json};
use tokio::time::{sleep, timeout};
use uuid::Uuid;
use ironflow_core::provider::{AgentProvider, LABEL_ROOT_RUN_ID};
use ironflow_core::providers::claude::ClaudeCodeProvider;
use ironflow_core::providers::record_replay::RecordReplayProvider;
use ironflow_engine::config::{DelayConfig, HumanInputConfig, ShellConfig, WorkflowOptions};
use ironflow_engine::context::{PARENT_RUN_ID_LABEL, WorkflowContext};
use ironflow_engine::engine::Engine;
use ironflow_engine::error::EngineError;
use ironflow_engine::executor::SubWorkflowOutcome;
use ironflow_engine::guard::WorkflowGuardConfig;
use ironflow_engine::handler::{HandlerFuture, TypedWorkflow, WorkflowHandler};
use ironflow_engine::plan::{ConditionResult, PlanOptions};
use ironflow_engine::signal::Signal;
use ironflow_engine::wake::RunWaker;
use ironflow_store::memory::InMemoryStore;
use ironflow_store::models::{
NewRun, Run, RunFilter, RunStatus, RunUpdate, Step, StepKind, StepStatus, StepUpdate,
TriggerKind,
};
use ironflow_store::store::{RunStore, Store};
const TEST_TIMEOUT: Duration = Duration::from_secs(10);
#[derive(Debug, Serialize, Deserialize, JsonSchema)]
struct GreetInput {
name: String,
shout: bool,
}
struct Greeter;
impl WorkflowHandler for Greeter {
fn name(&self) -> &str {
"greeter"
}
fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
Box::pin(async move {
if ctx.when("shouting", |i: &GreetInput| i.shout).await? {
ctx.shell("greet-loud", ShellConfig::new("echo HELLO"))
.await?;
} else {
ctx.shell("greet", ShellConfig::new("echo hello")).await?;
}
Ok(())
})
}
}
impl TypedWorkflow for Greeter {
type Input = GreetInput;
}
struct Host {
child_run_id: Arc<Mutex<Option<Uuid>>>,
}
impl WorkflowHandler for Host {
fn name(&self) -> &str {
"host"
}
fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
Box::pin(async move {
let child = ctx
.workflow(
&Greeter,
GreetInput {
name: "Élodie".to_string(),
shout: true,
},
)
.await?;
*self.child_run_id.lock().expect("lock") = Some(child.run_id());
Ok(())
})
}
}
fn provider() -> Arc<dyn AgentProvider> {
Arc::new(RecordReplayProvider::replay(
ClaudeCodeProvider::new(),
"/tmp/ironflow-fixtures",
))
}
fn engine(store: Arc<InMemoryStore>, child_run_id: Arc<Mutex<Option<Uuid>>>) -> Engine {
let store: Arc<dyn Store> = store;
let mut engine = Engine::new(store, provider());
engine.register(Greeter).expect("register child");
engine
.register(Host { child_run_id })
.expect("register parent");
engine
}
#[tokio::test]
async fn a_typed_sub_workflow_runs_with_its_input_and_reports_its_run_id() {
timeout(TEST_TIMEOUT, async {
let store = Arc::new(InMemoryStore::new());
let seen = Arc::new(Mutex::new(None));
let engine = engine(store.clone(), seen.clone());
let result = engine
.run_handler("host", TriggerKind::Manual, json!({}))
.await
.expect("parent completes");
assert_eq!(result.run.status.state, RunStatus::Completed);
let child = store
.list_runs(RunFilter::default(), 1, 10)
.await
.expect("list runs")
.items
.into_iter()
.find(|r| r.workflow_name == "greeter")
.expect("child run exists");
assert_eq!(*seen.lock().expect("lock"), Some(child.id));
assert_eq!(child.payload, json!({"name": "Élodie", "shout": true}));
let child_steps = store.list_steps(child.id).await.expect("child steps");
assert_eq!(child_steps.len(), 1);
assert_eq!(child_steps[0].name, "greet-loud");
})
.await
.expect("test timed out");
}
#[tokio::test]
async fn a_planned_sub_workflow_is_expanded_with_its_typed_input_and_a_nil_run_id() {
timeout(TEST_TIMEOUT, async {
let store = Arc::new(InMemoryStore::new());
let seen = Arc::new(Mutex::new(None));
let engine = engine(store.clone(), seen.clone());
let plan = engine
.plan_handler("host", json!({}), PlanOptions::default())
.await
.expect("plan built");
assert!(
!plan.truncated,
"plan stopped: {:?}",
plan.incomplete_reason
);
let names: Vec<&str> = plan.steps.iter().map(|s| s.name.as_str()).collect();
assert_eq!(names, vec!["greeter", "greet-loud"]);
assert_eq!(plan.steps[0].kind, StepKind::Workflow);
match plan.steps[1].condition.as_ref().expect("a condition") {
ConditionResult::Evaluated { expression, value } => {
assert_eq!(expression, "shouting");
assert!(*value);
}
other => panic!("expected an evaluated condition, got {other:?}"),
}
assert_eq!(*seen.lock().expect("lock"), Some(Uuid::nil()));
assert!(
store
.list_runs(RunFilter::default(), 1, 10)
.await
.expect("list runs")
.items
.is_empty(),
"planning must not create runs"
);
})
.await
.expect("test timed out");
}
const PARENT: &str = "parent";
const GRANDPARENT: &str = "grandparent";
const SIGNAL_KEY: &str = "release-42";
#[derive(Debug, Serialize, Deserialize, JsonSchema)]
struct NoInput {}
#[derive(Debug, Deserialize, JsonSchema)]
struct NameAnswer {
name: String,
}
#[derive(Debug, Serialize, Deserialize, JsonSchema)]
struct Deployed {
version: String,
}
impl Signal for Deployed {
const NAME: &'static str = "test.deployed";
}
type Seen = Arc<Mutex<Vec<String>>>;
fn seen(seen: &Seen) -> Vec<String> {
seen.lock().expect("seen lock").clone()
}
struct Asker {
seen: Seen,
}
impl WorkflowHandler for Asker {
fn name(&self) -> &str {
"asker"
}
fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
Box::pin(async move {
let answer: NameAnswer = ctx
.human_input("ask-name", HumanInputConfig::new("Who ships it?"))
.await?;
self.seen.lock().expect("seen lock").push(answer.name);
ctx.shell("after-input", ShellConfig::new("echo answered"))
.await?;
Ok(())
})
}
}
impl TypedWorkflow for Asker {
type Input = NoInput;
}
struct Waiter {
seen: Seen,
}
impl WorkflowHandler for Waiter {
fn name(&self) -> &str {
"waiter"
}
fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
Box::pin(async move {
let deployed = ctx
.wait_for_signal::<Deployed>("wait-deploy", SIGNAL_KEY, Duration::from_secs(3600))
.await?;
let seen_value = match deployed {
Some(deployed) => deployed.version,
None => "timed out".to_string(),
};
self.seen.lock().expect("seen lock").push(seen_value);
Ok(())
})
}
}
impl TypedWorkflow for Waiter {
type Input = NoInput;
}
struct Sleeper {
seen: Seen,
}
impl WorkflowHandler for Sleeper {
fn name(&self) -> &str {
"sleeper"
}
fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
Box::pin(async move {
ctx.delay("pause", DelayConfig::from_secs(300)).await?;
self.seen
.lock()
.expect("seen lock")
.push("slept".to_string());
Ok(())
})
}
}
impl TypedWorkflow for Sleeper {
type Input = NoInput;
}
#[derive(Clone, Copy)]
enum Suspends {
HumanInput,
Signal,
Delay,
}
struct Parent {
suspends: Suspends,
seen: Seen,
}
impl WorkflowHandler for Parent {
fn name(&self) -> &str {
PARENT
}
fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
Box::pin(async move {
ctx.workflow(
&Greeter,
GreetInput {
name: "sibling".to_string(),
shout: false,
},
)
.await?;
let seen = self.seen.clone();
match self.suspends {
Suspends::HumanInput => ctx.workflow(&Asker { seen }, NoInput {}).await?,
Suspends::Signal => ctx.workflow(&Waiter { seen }, NoInput {}).await?,
Suspends::Delay => ctx.workflow(&Sleeper { seen }, NoInput {}).await?,
};
ctx.shell("after-child", ShellConfig::new("echo done"))
.await?;
Ok(())
})
}
}
impl TypedWorkflow for Parent {
type Input = NoInput;
}
struct Grandparent {
suspends: Suspends,
seen: Seen,
}
impl WorkflowHandler for Grandparent {
fn name(&self) -> &str {
GRANDPARENT
}
fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
Box::pin(async move {
let parent = Parent {
suspends: self.suspends,
seen: self.seen.clone(),
};
ctx.workflow(&parent, NoInput {}).await?;
Ok(())
})
}
}
fn new_engine(store: &Arc<InMemoryStore>) -> Engine {
let store: Arc<dyn Store> = store.clone();
Engine::new(store, provider())
}
fn build_chain(mut engine: Engine, suspends: Suspends) -> (Arc<Engine>, Seen) {
let seen = Seen::default();
engine.register(Greeter).expect("register greeter");
engine
.register(Asker { seen: seen.clone() })
.expect("register asker");
engine
.register(Waiter { seen: seen.clone() })
.expect("register waiter");
engine
.register(Sleeper { seen: seen.clone() })
.expect("register sleeper");
engine
.register(Parent {
suspends,
seen: seen.clone(),
})
.expect("register parent");
engine
.register(Grandparent {
suspends,
seen: seen.clone(),
})
.expect("register grandparent");
(Arc::new(engine), seen)
}
async fn start_chain(engine: &Engine, root: &str) -> Run {
engine
.run_handler(root, TriggerKind::Manual, json!({}))
.await
.expect("the chain suspends")
.run
}
async fn load_run(store: &InMemoryStore, run_id: Uuid) -> Run {
store
.get_run(run_id)
.await
.expect("get run")
.expect("run exists")
}
async fn run_of(store: &InMemoryStore, workflow: &str) -> Run {
let filter = RunFilter {
workflow_name: Some(workflow.to_string()),
..RunFilter::default()
};
let mut runs = store
.list_runs(filter, 1, 50)
.await
.expect("list runs")
.items;
runs.retain(|r| r.workflow_name == workflow);
assert_eq!(runs.len(), 1, "expected exactly one {workflow} run");
runs.remove(0)
}
async fn workflow_step(store: &InMemoryStore, run_id: Uuid, child: &str) -> Step {
store
.list_steps(run_id)
.await
.expect("list steps")
.into_iter()
.find(|s| s.kind == StepKind::Workflow && s.name == child)
.expect("the workflow step exists")
}
async fn open_input_step(store: &InMemoryStore, run_id: Uuid) -> Step {
let step = store
.list_steps(run_id)
.await
.expect("list steps")
.into_iter()
.find(|s| s.kind == StepKind::HumanInput)
.expect("a human input step");
assert_eq!(step.status.state, StepStatus::AwaitingApproval);
step
}
async fn answer(store: &InMemoryStore, run_id: Uuid, name: &str) {
let step = open_input_step(store, run_id).await;
store
.update_step(
step.id,
StepUpdate {
status: Some(StepStatus::Completed),
output: Some(json!({ "name": name })),
completed_at: Some(Utc::now()),
clear_approval_deadline: true,
..StepUpdate::default()
},
)
.await
.expect("store the answer");
store
.update_run_status(run_id, RunStatus::Running)
.await
.expect("mark running");
}
async fn reject(store: &InMemoryStore, run_id: Uuid, reason: &str) {
let step = open_input_step(store, run_id).await;
store
.update_step(
step.id,
StepUpdate {
status: Some(StepStatus::Rejected),
error: Some(reason.to_string()),
completed_at: Some(Utc::now()),
clear_approval_deadline: true,
..StepUpdate::default()
},
)
.await
.expect("reject the input");
store
.update_run_status(run_id, RunStatus::Running)
.await
.expect("mark running");
}
async fn wait_for_status(store: &InMemoryStore, run_id: Uuid, status: RunStatus) {
timeout(TEST_TIMEOUT, async {
loop {
if load_run(store, run_id).await.status.state == status {
return;
}
sleep(Duration::from_millis(20)).await;
}
})
.await
.expect("run never reached the expected status");
}
#[tokio::test]
async fn sub_workflow_human_input_suspends_the_parent_and_resumes_through_the_root() {
timeout(TEST_TIMEOUT, async {
let store = Arc::new(InMemoryStore::new());
let (engine, seen_names) = build_chain(new_engine(&store), Suspends::HumanInput);
let parent = start_chain(&engine, PARENT).await;
assert_eq!(parent.status.state, RunStatus::AwaitingApproval);
assert!(parent.scheduled_at.is_none());
let child = run_of(&store, "asker").await;
assert_eq!(child.status.state, RunStatus::AwaitingApproval);
assert_eq!(child.trigger, TriggerKind::Workflow);
let step = workflow_step(&store, parent.id, "asker").await;
assert_eq!(
step.status.state,
StepStatus::Running,
"the step stays open while its child is suspended"
);
assert_eq!(step.output, Some(json!({ "child_run_id": child.id })));
assert!(seen(&seen_names).is_empty());
answer(&store, child.id, "Ada").await;
let result = engine
.resume_run(child.id)
.await
.expect("the chain resumes");
assert_eq!(result.run.id, parent.id, "the root run is the one resumed");
assert_eq!(result.run.status.state, RunStatus::Completed);
assert_eq!(
run_of(&store, "asker").await.id,
child.id,
"the same child run is re-entered, never a new one"
);
assert_eq!(
load_run(&store, child.id).await.status.state,
RunStatus::Completed
);
let step = workflow_step(&store, parent.id, "asker").await;
assert_eq!(step.status.state, StepStatus::Completed);
let output: Value = step.output.expect("the step has an output");
assert_eq!(output["run_id"], json!(child.id));
assert_eq!(output["status"], json!("completed"));
assert_eq!(seen(&seen_names), vec!["Ada".to_string()]);
let parent_steps = store.list_steps(parent.id).await.expect("list steps");
let names: Vec<&str> = parent_steps.iter().map(|s| s.name.as_str()).collect();
assert_eq!(names, vec!["greeter", "asker", "after-child"]);
})
.await
.expect("test timed out");
}
#[tokio::test]
async fn sub_workflow_completed_sibling_is_not_run_again_on_resume() {
timeout(TEST_TIMEOUT, async {
let store = Arc::new(InMemoryStore::new());
let (engine, _seen) = build_chain(new_engine(&store), Suspends::HumanInput);
let parent = start_chain(&engine, PARENT).await;
let sibling = run_of(&store, "greeter").await;
assert_eq!(sibling.status.state, RunStatus::Completed);
let child = run_of(&store, "asker").await;
answer(&store, child.id, "Ada").await;
engine
.resume_run(child.id)
.await
.expect("the chain resumes");
assert_eq!(
run_of(&store, "greeter").await.id,
sibling.id,
"a completed child is replayed, never run again"
);
let sibling_steps = store.list_steps(sibling.id).await.expect("list steps");
assert_eq!(sibling_steps.len(), 1);
let step = workflow_step(&store, parent.id, "greeter").await;
assert_eq!(step.status.state, StepStatus::Completed);
})
.await
.expect("test timed out");
}
#[tokio::test]
async fn sub_workflow_suspended_child_is_listed_by_its_chain_labels() {
timeout(TEST_TIMEOUT, async {
let store = Arc::new(InMemoryStore::new());
let (engine, _seen) = build_chain(new_engine(&store), Suspends::HumanInput);
let parent = start_chain(&engine, PARENT).await;
let waiting = store
.list_runs(
RunFilter {
labels: Some(HashMap::from([(
PARENT_RUN_ID_LABEL.to_string(),
parent.id.to_string(),
)])),
status: Some(RunStatus::AwaitingApproval),
..RunFilter::default()
},
1,
50,
)
.await
.expect("list runs")
.items;
assert_eq!(waiting.len(), 1);
assert_eq!(waiting[0].workflow_name, "asker");
let chain = store
.list_runs(
RunFilter {
labels: Some(HashMap::from([(
LABEL_ROOT_RUN_ID.to_string(),
parent.id.to_string(),
)])),
..RunFilter::default()
},
1,
50,
)
.await
.expect("list runs")
.items;
let mut names: Vec<&str> = chain.iter().map(|r| r.workflow_name.as_str()).collect();
names.sort_unstable();
assert_eq!(names, vec!["asker", "greeter"]);
})
.await
.expect("test timed out");
}
#[tokio::test]
async fn sub_workflow_child_executed_by_a_worker_resumes_its_root() {
timeout(TEST_TIMEOUT, async {
let store = Arc::new(InMemoryStore::new());
let (engine, seen_names) = build_chain(new_engine(&store), Suspends::HumanInput);
let parent = start_chain(&engine, PARENT).await;
let child = run_of(&store, "asker").await;
answer(&store, child.id, "Grace").await;
let result = engine
.execute_handler_run(child.id)
.await
.expect("the chain resumes");
assert_eq!(result.run.id, parent.id);
assert_eq!(result.run.status.state, RunStatus::Completed);
assert_eq!(
load_run(&store, child.id).await.status.state,
RunStatus::Completed
);
assert_eq!(seen(&seen_names), vec!["Grace".to_string()]);
})
.await
.expect("test timed out");
}
#[tokio::test]
async fn sub_workflow_signal_wait_suspends_the_chain_until_the_signal_arrives() {
timeout(TEST_TIMEOUT, async {
let store = Arc::new(InMemoryStore::new());
let (engine, seen_versions) = build_chain(new_engine(&store), Suspends::Signal);
let parent = start_chain(&engine, PARENT).await;
assert_eq!(parent.status.state, RunStatus::Sleeping);
assert!(
parent.scheduled_at.is_none(),
"only the waiting child owns the wake-up"
);
let child = run_of(&store, "waiter").await;
assert_eq!(child.status.state, RunStatus::Sleeping);
assert!(child.scheduled_at.is_some(), "the child arms the deadline");
let delivery = engine
.send_signal(
&Deployed {
version: "1.2.3".to_string(),
},
SIGNAL_KEY,
None,
)
.await
.expect("deliver");
assert_eq!(delivery.resumed.len(), 1);
assert_eq!(delivery.resumed[0].run_id, child.id);
wait_for_status(&store, parent.id, RunStatus::Completed).await;
assert_eq!(
load_run(&store, child.id).await.status.state,
RunStatus::Completed
);
assert_eq!(seen(&seen_versions), vec!["1.2.3".to_string()]);
assert_eq!(run_of(&store, "waiter").await.id, child.id);
})
.await
.expect("test timed out");
}
#[tokio::test]
async fn sub_workflow_delay_wakes_the_chain_through_the_waker() {
timeout(TEST_TIMEOUT, async {
let store = Arc::new(InMemoryStore::new());
let (engine, seen_wakes) = build_chain(new_engine(&store), Suspends::Delay);
let parent = start_chain(&engine, PARENT).await;
assert_eq!(parent.status.state, RunStatus::Sleeping);
assert!(parent.scheduled_at.is_none());
let child = run_of(&store, "sleeper").await;
assert_eq!(child.status.state, RunStatus::Sleeping);
assert!(child.scheduled_at.is_some());
store
.update_run(
child.id,
RunUpdate {
scheduled_at: Some(Utc::now() - TimeDelta::seconds(1)),
..RunUpdate::default()
},
)
.await
.expect("move the wake-up");
let woken = RunWaker::new(engine.clone()).tick().await.expect("tick");
assert_eq!(
woken.iter().map(|r| r.id).collect::<Vec<_>>(),
vec![child.id],
"the parent is never woken on its own"
);
wait_for_status(&store, parent.id, RunStatus::Completed).await;
assert_eq!(
load_run(&store, child.id).await.status.state,
RunStatus::Completed
);
assert_eq!(seen(&seen_wakes), vec!["slept".to_string()]);
})
.await
.expect("test timed out");
}
#[tokio::test]
async fn sub_workflow_rejected_human_input_fails_the_child_and_the_parent() {
timeout(TEST_TIMEOUT, async {
let store = Arc::new(InMemoryStore::new());
let (engine, seen_names) = build_chain(new_engine(&store), Suspends::HumanInput);
let parent = start_chain(&engine, PARENT).await;
let child = run_of(&store, "asker").await;
reject(&store, child.id, "nobody").await;
let err = engine
.resume_run(child.id)
.await
.expect_err("a rejected input fails the chain");
assert!(
matches!(err, EngineError::HumanInputRejected { .. }),
"got {err}"
);
assert_eq!(
load_run(&store, child.id).await.status.state,
RunStatus::Failed
);
assert_eq!(
load_run(&store, parent.id).await.status.state,
RunStatus::Failed
);
let step = workflow_step(&store, parent.id, "asker").await;
assert_eq!(step.status.state, StepStatus::Failed);
assert!(seen(&seen_names).is_empty());
})
.await
.expect("test timed out");
}
#[tokio::test]
async fn sub_workflow_child_of_a_cancelled_root_cannot_resume() {
timeout(TEST_TIMEOUT, async {
let store = Arc::new(InMemoryStore::new());
let (engine, seen_names) = build_chain(new_engine(&store), Suspends::HumanInput);
let parent = start_chain(&engine, PARENT).await;
let child = run_of(&store, "asker").await;
store
.update_run_status(parent.id, RunStatus::Cancelled)
.await
.expect("cancel the root");
answer(&store, child.id, "Ada").await;
let err = engine
.resume_run(child.id)
.await
.expect_err("a cancelled root cannot take its child back");
assert!(matches!(err, EngineError::InvalidWorkflow(_)), "got {err}");
assert_eq!(
load_run(&store, child.id).await.status.state,
RunStatus::Failed
);
assert_eq!(
load_run(&store, parent.id).await.status.state,
RunStatus::Cancelled
);
assert!(seen(&seen_names).is_empty());
})
.await
.expect("test timed out");
}
#[tokio::test]
async fn sub_workflow_nested_grandchild_suspends_and_resumes_the_whole_chain() {
timeout(TEST_TIMEOUT, async {
let store = Arc::new(InMemoryStore::new());
let (engine, seen_names) = build_chain(new_engine(&store), Suspends::HumanInput);
let root = start_chain(&engine, GRANDPARENT).await;
assert_eq!(root.status.state, RunStatus::AwaitingApproval);
let parent = run_of(&store, PARENT).await;
assert_eq!(parent.status.state, RunStatus::AwaitingApproval);
assert!(
parent.scheduled_at.is_none(),
"an intermediate run never owns a wake-up"
);
assert_eq!(
parent.labels.get(PARENT_RUN_ID_LABEL),
Some(&root.id.to_string())
);
let grandchild = run_of(&store, "asker").await;
assert_eq!(grandchild.status.state, RunStatus::AwaitingApproval);
assert_eq!(
grandchild.labels.get(PARENT_RUN_ID_LABEL),
Some(&parent.id.to_string())
);
assert_eq!(
grandchild.labels.get(LABEL_ROOT_RUN_ID),
Some(&root.id.to_string())
);
answer(&store, grandchild.id, "Ada").await;
let result = engine
.resume_run(grandchild.id)
.await
.expect("the chain resumes");
assert_eq!(result.run.id, root.id);
assert_eq!(result.run.status.state, RunStatus::Completed);
for run_id in [parent.id, grandchild.id] {
assert_eq!(
load_run(&store, run_id).await.status.state,
RunStatus::Completed
);
}
assert_eq!(run_of(&store, PARENT).await.id, parent.id);
assert_eq!(run_of(&store, "asker").await.id, grandchild.id);
assert_eq!(
run_of(&store, "greeter").await.status.state,
RunStatus::Completed
);
assert_eq!(seen(&seen_names), vec!["Ada".to_string()]);
})
.await
.expect("test timed out");
}
#[tokio::test]
async fn sub_workflow_resume_stays_within_the_guard_fan_out() {
timeout(TEST_TIMEOUT, async {
let store = Arc::new(InMemoryStore::new());
let engine =
new_engine(&store).with_guard_config(WorkflowGuardConfig::new().with_max_fan_out(2));
let (engine, seen_names) = build_chain(engine, Suspends::HumanInput);
start_chain(&engine, PARENT).await;
let child = run_of(&store, "asker").await;
answer(&store, child.id, "Ada").await;
let result = engine
.resume_run(child.id)
.await
.expect("the guard lets the chain resume");
assert_eq!(result.run.status.state, RunStatus::Completed);
assert_eq!(seen(&seen_names), vec!["Ada".to_string()]);
})
.await
.expect("test timed out");
}
const ISSUE_KEY: &str = "issue:12";
type Outcomes = Arc<Mutex<Vec<String>>>;
struct Dispatcher {
outcomes: Outcomes,
}
impl WorkflowHandler for Dispatcher {
fn name(&self) -> &str {
"dispatcher"
}
fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
Box::pin(async move {
let input = GreetInput {
name: "Ada".to_string(),
shout: false,
};
let options = WorkflowOptions::new().concurrency_key(ISSUE_KEY);
let seen = match ctx.workflow_with(&Greeter, input, options).await? {
SubWorkflowOutcome::Completed(child) => format!("completed:{}", child.run_id()),
SubWorkflowOutcome::Conflict(c) => format!("conflict:{}:{}", c.key(), c.run_id()),
};
self.outcomes.lock().expect("outcomes lock").push(seen);
let _confirmed: NameAnswer = ctx
.human_input("confirm", HumanInputConfig::new("Who confirms?"))
.await?;
ctx.shell("after-confirm", ShellConfig::new("echo confirmed"))
.await?;
Ok(())
})
}
}
fn dispatcher_engine(store: &Arc<InMemoryStore>) -> (Engine, Outcomes) {
let outcomes = Outcomes::default();
let mut engine = new_engine(store);
engine.register(Greeter).expect("register greeter");
engine
.register(Dispatcher {
outcomes: outcomes.clone(),
})
.expect("register dispatcher");
(engine, outcomes)
}
fn outcomes(outcomes: &Outcomes) -> Vec<String> {
outcomes.lock().expect("outcomes lock").clone()
}
async fn runs_of(store: &InMemoryStore, workflow: &str) -> Vec<Run> {
let filter = RunFilter {
workflow_name: Some(workflow.to_string()),
..RunFilter::default()
};
let mut runs = store
.list_runs(filter, 1, 50)
.await
.expect("list runs")
.items;
runs.retain(|r| r.workflow_name == workflow);
runs
}
async fn create_blocker(store: &InMemoryStore) -> Run {
store
.create_run(NewRun {
workflow_name: "blocker".to_string(),
trigger: TriggerKind::Manual,
payload: json!({}),
max_retries: 0,
handler_version: None,
labels: HashMap::new(),
scheduled_at: None,
created_by: None,
idempotency_key: None,
concurrency_key: Some(ISSUE_KEY.to_string()),
max_cost_usd: None,
})
.await
.expect("create the blocker")
.into_run()
}
#[tokio::test]
async fn sub_workflow_concurrency_conflict_completes_the_step_without_a_child() {
timeout(TEST_TIMEOUT, async {
let store = Arc::new(InMemoryStore::new());
let (engine, seen_outcomes) = dispatcher_engine(&store);
let blocker = create_blocker(&store).await;
let dispatcher = engine
.run_handler("dispatcher", TriggerKind::Manual, json!({}))
.await
.expect("the dispatcher suspends on its human input")
.run;
assert_eq!(
dispatcher.status.state,
RunStatus::AwaitingApproval,
"a conflict does not fail the parent"
);
assert!(
runs_of(&store, "greeter").await.is_empty(),
"no child run is created on a conflict"
);
assert_eq!(
outcomes(&seen_outcomes),
vec![format!("conflict:{ISSUE_KEY}:{}", blocker.id)]
);
let step = workflow_step(&store, dispatcher.id, "greeter").await;
assert_eq!(step.status.state, StepStatus::Completed);
assert_eq!(
step.output,
Some(json!({
"concurrency_conflict": { "key": ISSUE_KEY, "run_id": blocker.id }
}))
);
assert_eq!(step.duration_ms, 0);
})
.await
.expect("test timed out");
}
#[tokio::test]
async fn sub_workflow_concurrency_conflict_is_replayed_without_a_new_child() {
timeout(TEST_TIMEOUT, async {
let store = Arc::new(InMemoryStore::new());
let (engine, seen_outcomes) = dispatcher_engine(&store);
let blocker = create_blocker(&store).await;
let dispatcher = engine
.run_handler("dispatcher", TriggerKind::Manual, json!({}))
.await
.expect("the dispatcher suspends on its human input")
.run;
assert_eq!(dispatcher.status.state, RunStatus::AwaitingApproval);
store
.update_run_status(blocker.id, RunStatus::Cancelled)
.await
.expect("cancel the blocker");
answer(&store, dispatcher.id, "Grace").await;
let result = engine
.resume_run(dispatcher.id)
.await
.expect("the dispatcher resumes");
assert_eq!(result.run.status.state, RunStatus::Completed);
assert!(
runs_of(&store, "greeter").await.is_empty(),
"the replay does not create a child"
);
let conflict = format!("conflict:{ISSUE_KEY}:{}", blocker.id);
assert_eq!(outcomes(&seen_outcomes), vec![conflict.clone(), conflict]);
let steps = store.list_steps(dispatcher.id).await.expect("list steps");
let names: Vec<&str> = steps.iter().map(|s| s.name.as_str()).collect();
assert_eq!(names, vec!["greeter", "confirm", "after-confirm"]);
assert_eq!(
steps[0].output,
Some(json!({
"concurrency_conflict": { "key": ISSUE_KEY, "run_id": blocker.id }
})),
"the recorded conflict is left as it was"
);
})
.await
.expect("test timed out");
}
#[tokio::test]
async fn sub_workflow_with_concurrency_key_creates_a_child_holding_the_key() {
timeout(TEST_TIMEOUT, async {
let store = Arc::new(InMemoryStore::new());
let (engine, seen_outcomes) = dispatcher_engine(&store);
let dispatcher = engine
.run_handler("dispatcher", TriggerKind::Manual, json!({}))
.await
.expect("the dispatcher suspends on its human input")
.run;
assert_eq!(dispatcher.status.state, RunStatus::AwaitingApproval);
let child = run_of(&store, "greeter").await;
assert_eq!(child.status.state, RunStatus::Completed);
assert_eq!(child.concurrency_key.as_deref(), Some(ISSUE_KEY));
assert_eq!(
outcomes(&seen_outcomes),
vec![format!("completed:{}", child.id)]
);
let step = workflow_step(&store, dispatcher.id, "greeter").await;
assert_eq!(step.status.state, StepStatus::Completed);
assert_eq!(
step.input.expect("the step records its config")["concurrency_key"],
json!(ISSUE_KEY)
);
let next = create_blocker(&store).await;
assert_eq!(next.concurrency_key.as_deref(), Some(ISSUE_KEY));
})
.await
.expect("test timed out");
}
struct Failer;
impl WorkflowHandler for Failer {
fn name(&self) -> &str {
"failer"
}
fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
Box::pin(async move {
ctx.shell("boom", ShellConfig::new("exit 1")).await?;
Ok(())
})
}
}
impl TypedWorkflow for Failer {
type Input = NoInput;
}
struct Tolerant {
seen: Seen,
}
impl WorkflowHandler for Tolerant {
fn name(&self) -> &str {
"tolerant"
}
fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
Box::pin(async move {
let failed = match ctx
.workflow_with(&Failer, NoInput {}, WorkflowOptions::new().allow_failure())
.await?
{
SubWorkflowOutcome::Completed(output) => output,
SubWorkflowOutcome::Conflict(_) => {
return Err(EngineError::InvalidWorkflow(
"no concurrency key was set".to_string(),
));
}
};
if failed.status() != RunStatus::Failed || failed.error().is_none() {
return Err(EngineError::InvalidWorkflow(format!(
"the failed child was not reported: {:?} {:?}",
failed.status(),
failed.error()
)));
}
ctx.workflow(
&Asker {
seen: self.seen.clone(),
},
NoInput {},
)
.await?;
ctx.shell("after-child", ShellConfig::new("echo done"))
.await?;
Ok(())
})
}
}
impl TypedWorkflow for Tolerant {
type Input = NoInput;
}
struct Strict;
impl WorkflowHandler for Strict {
fn name(&self) -> &str {
"strict"
}
fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
Box::pin(async move {
ctx.workflow(&Failer, NoInput {}).await?;
Ok(())
})
}
}
impl TypedWorkflow for Strict {
type Input = NoInput;
}
struct TolerantAsker {
seen: Seen,
}
impl WorkflowHandler for TolerantAsker {
fn name(&self) -> &str {
"tolerant-asker"
}
fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
Box::pin(async move {
ctx.workflow_with(
&Asker {
seen: self.seen.clone(),
},
NoInput {},
WorkflowOptions::new().allow_failure(),
)
.await?;
Ok(())
})
}
}
impl TypedWorkflow for TolerantAsker {
type Input = NoInput;
}
fn allow_failure_engine(store: &Arc<InMemoryStore>) -> (Arc<Engine>, Seen) {
let seen = Seen::default();
let mut engine = new_engine(store);
engine.register(Failer).expect("register failer");
engine
.register(Asker { seen: seen.clone() })
.expect("register asker");
engine
.register(Tolerant { seen: seen.clone() })
.expect("register tolerant");
engine.register(Strict).expect("register strict");
engine
.register(TolerantAsker { seen: seen.clone() })
.expect("register tolerant asker");
(Arc::new(engine), seen)
}
async fn count_runs(store: &InMemoryStore, workflow: &str) -> usize {
let filter = RunFilter {
workflow_name: Some(workflow.to_string()),
..RunFilter::default()
};
store
.list_runs(filter, 1, 50)
.await
.expect("list runs")
.items
.iter()
.filter(|r| r.workflow_name == workflow)
.count()
}
#[tokio::test]
async fn sub_workflow_allow_failure_tolerates_a_failed_child_then_resumes_without_a_new_child() {
timeout(TEST_TIMEOUT, async {
let store = Arc::new(InMemoryStore::new());
let (engine, seen_names) = allow_failure_engine(&store);
let parent = start_chain(&engine, "tolerant").await;
assert_eq!(parent.status.state, RunStatus::AwaitingApproval);
let failer_step = workflow_step(&store, parent.id, "failer").await;
assert_eq!(failer_step.status.state, StepStatus::Completed);
assert_eq!(count_runs(&store, "failer").await, 1);
let asker = run_of(&store, "asker").await;
answer(&store, asker.id, "Ada").await;
let result = engine
.resume_run(asker.id)
.await
.expect("the chain resumes");
assert_eq!(result.run.id, parent.id);
assert_eq!(result.run.status.state, RunStatus::Warning);
assert_eq!(
count_runs(&store, "failer").await,
1,
"the resume replays the step and creates no new child"
);
let replayed = workflow_step(&store, parent.id, "failer").await;
assert_eq!(replayed.id, failer_step.id);
assert_eq!(replayed.status.state, StepStatus::Completed);
let failer = run_of(&store, "failer").await;
assert_eq!(failer.status.state, RunStatus::Failed);
assert!(failer.error.is_some(), "the child error is recorded");
let output: Value = replayed.output.expect("the step has an output");
assert_eq!(output["run_id"], json!(failer.id));
assert_eq!(output["status"], json!("failed"));
assert!(output["error"].is_string());
assert_eq!(seen(&seen_names), vec!["Ada".to_string()]);
})
.await
.expect("test timed out");
}
#[tokio::test]
async fn sub_workflow_without_allow_failure_a_failed_child_still_fails_the_parent() {
timeout(TEST_TIMEOUT, async {
let store = Arc::new(InMemoryStore::new());
let (engine, _seen) = allow_failure_engine(&store);
let outcome = engine
.run_handler("strict", TriggerKind::Manual, json!({}))
.await;
drop(outcome);
let parent = run_of(&store, "strict").await;
assert_eq!(parent.status.state, RunStatus::Failed);
let step = workflow_step(&store, parent.id, "failer").await;
assert_eq!(step.status.state, StepStatus::Failed);
assert_eq!(
run_of(&store, "failer").await.status.state,
RunStatus::Failed
);
})
.await
.expect("test timed out");
}
#[tokio::test]
async fn sub_workflow_allow_failure_does_not_tolerate_a_suspension() {
timeout(TEST_TIMEOUT, async {
let store = Arc::new(InMemoryStore::new());
let (engine, seen_names) = allow_failure_engine(&store);
let parent = start_chain(&engine, "tolerant-asker").await;
assert_eq!(parent.status.state, RunStatus::AwaitingApproval);
let step = workflow_step(&store, parent.id, "asker").await;
assert_eq!(step.status.state, StepStatus::Running);
let child = run_of(&store, "asker").await;
answer(&store, child.id, "Ada").await;
let result = engine
.resume_run(child.id)
.await
.expect("the chain resumes");
assert_eq!(result.run.status.state, RunStatus::Completed);
assert_eq!(
workflow_step(&store, parent.id, "asker").await.status.state,
StepStatus::Completed
);
assert_eq!(seen(&seen_names), vec!["Ada".to_string()]);
})
.await
.expect("test timed out");
}