use super::model::{OnError, Step, Workflow};
use super::template::{self, Data};
use crate::state::now_ms;
use serde::{Deserialize, Serialize};
use serde_json::{Map, Value, json};
use std::collections::BTreeMap;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Default)]
#[serde(rename_all = "snake_case")]
pub enum StepStatus {
#[default]
Pending,
Running,
Done,
Failed,
Skipped,
Cancelled,
Timeout,
Suspended,
Pruned,
}
impl StepStatus {
pub fn is_terminal(self) -> bool {
matches!(
self,
StepStatus::Done
| StepStatus::Failed
| StepStatus::Skipped
| StepStatus::Pruned
| StepStatus::Cancelled
| StepStatus::Timeout
)
}
pub fn is_satisfied(self) -> bool {
matches!(self, StepStatus::Done | StepStatus::Skipped)
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, Default)]
pub struct StepState {
#[serde(default)]
pub status: StepStatus,
#[serde(default)]
pub attempt: u32,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub started: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub finished: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub output: Option<Value>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub error: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub wait: Option<Value>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub cache_key: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub worker: Option<String>,
#[serde(default, skip_serializing_if = "std::ops::Not::not")]
pub forced: bool,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Default)]
#[serde(rename_all = "snake_case")]
pub enum RunStatus {
#[default]
Pending,
Running,
Suspended,
Paused,
Completed,
Failed,
Refused,
Cancelled,
Stalled,
}
impl RunStatus {
pub fn is_terminal(self) -> bool {
matches!(
self,
RunStatus::Completed
| RunStatus::Failed
| RunStatus::Refused
| RunStatus::Cancelled
| RunStatus::Stalled
)
}
pub fn as_str(self) -> &'static str {
match self {
RunStatus::Pending => "pending",
RunStatus::Running => "running",
RunStatus::Suspended => "suspended",
RunStatus::Paused => "paused",
RunStatus::Completed => "completed",
RunStatus::Failed => "failed",
RunStatus::Refused => "refused",
RunStatus::Cancelled => "cancelled",
RunStatus::Stalled => "stalled",
}
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, Default)]
pub struct Start {
pub node: String,
#[serde(default)]
pub payload: Value,
#[serde(default)]
pub ts: u64,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct RunState {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub break_before: Option<String>,
pub id: String,
pub workflow: String,
pub workflow_hash: String,
#[serde(default)]
pub inputs: Value,
#[serde(default)]
pub status: RunStatus,
#[serde(default)]
pub start: Start,
#[serde(default)]
pub steps: BTreeMap<String, StepState>,
#[serde(default)]
pub vars: Map<String, Value>,
#[serde(default)]
pub tokens: u64,
#[serde(default)]
pub steps_run: u32,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub output: Option<Value>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub error: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub task: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub principal: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub conversation: Option<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub children: Vec<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub parent: Option<Value>,
#[serde(default)]
pub msg_depth: u32,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub key: Option<String>,
#[serde(default)]
pub attempt: u32,
#[serde(default)]
pub created: u64,
#[serde(default)]
pub updated: u64,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub finished: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub deadline_ms: Option<u64>,
#[serde(default = "default_durable")]
pub durable: bool,
#[serde(skip)]
pub dirty: bool,
}
fn default_durable() -> bool {
true
}
impl RunState {
pub fn new(id: &str, wf: &Workflow, start: Start, inputs: Value) -> RunState {
let now = now_ms();
let mut steps = BTreeMap::new();
for s in wf.steps.keys() {
steps.insert(s.clone(), StepState::default());
}
for s in wf.start_steps() {
let st = steps.get_mut(&s.id).expect("present");
if s.id == start.node {
st.status = StepStatus::Done;
st.output = Some(start.payload.clone());
st.started = Some(now);
st.finished = Some(now);
st.attempt = 1;
} else {
st.status = StepStatus::Skipped;
}
}
RunState {
break_before: None,
id: id.to_string(),
workflow: wf.name.clone(),
workflow_hash: wf.hash.clone(),
inputs,
status: RunStatus::Running,
start,
steps,
vars: Map::new(),
tokens: 0,
steps_run: 0,
output: None,
error: None,
task: None,
principal: None,
conversation: None,
children: Vec::new(),
parent: None,
msg_depth: 0,
key: None,
attempt: 1,
created: now,
updated: now,
finished: None,
deadline_ms: wf.limits.deadline_ms.map(|d| now + d),
durable: wf.durable.unwrap_or(true),
dirty: true,
}
}
pub fn touch(&mut self) {
self.updated = now_ms();
self.dirty = true;
}
pub fn step(&self, id: &str) -> Option<&StepState> {
self.steps.get(id)
}
pub fn data(&self, env: Value, memory: Value) -> Data {
let mut d = Data::new();
d.insert("inputs".into(), self.inputs.clone());
d.insert(
"run".into(),
json!({"id": self.id, "workflow": self.workflow, "start": self.start, "principal": self.principal, "task": self.task, "attempt": self.attempt, "status": self.status}),
);
d.insert(
"steps".into(),
Value::Object(
self.steps
.iter()
.map(|(k, s)| (k.clone(), json!({"status": s.status, "output": s.output, "error": s.error, "attempt": s.attempt})))
.collect(),
),
);
d.insert("vars".into(), Value::Object(self.vars.clone()));
d.insert("env".into(), env);
d.insert("memory".into(), memory);
d
}
pub fn begin_step(&mut self, id: &str) -> u32 {
let attempt = {
let st = self.steps.entry(id.to_string()).or_default();
st.status = StepStatus::Running;
st.attempt += 1;
st.started = Some(now_ms());
st.finished = None;
st.error = None;
st.wait = None;
st.cache_key = None;
st.forced = false;
st.attempt
};
if !self.status.is_terminal() {
self.status = RunStatus::Running;
}
self.touch();
attempt
}
pub fn end_step(
&mut self,
id: &str,
status: StepStatus,
output: Option<Value>,
error: Option<String>,
) {
let st = self.steps.entry(id.to_string()).or_default();
st.status = status;
st.finished = Some(now_ms());
st.output = output;
st.error = error;
st.wait = None;
st.cache_key = None;
st.worker = None;
self.steps_run += 1;
self.touch();
}
pub fn suspend_step(&mut self, id: &str, wait: Value) {
let st = self.steps.entry(id.to_string()).or_default();
st.status = StepStatus::Suspended;
st.wait = Some(wait);
self.touch();
}
pub fn finish(&mut self, status: RunStatus, output: Option<Value>, error: Option<String>) {
self.status = status;
self.output = output;
self.error = error;
self.finished = Some(now_ms());
for st in self.steps.values_mut() {
if !st.status.is_terminal() {
st.status = StepStatus::Cancelled;
st.finished = Some(now_ms());
}
}
self.touch();
}
pub fn write_var(&mut self, key: &str, value: Value, mode: &str) {
let cur = self.vars.remove(key);
let next = match (mode, cur) {
("append", Some(Value::Array(mut a))) => {
match value {
Value::Array(more) => a.extend(more),
other => a.push(other),
}
Value::Array(a)
}
("append", Some(other)) => json!([other, value]),
("append", None) => match value {
Value::Array(a) => Value::Array(a),
other => json!([other]),
},
("merge", Some(Value::Object(mut o))) => {
if let Value::Object(more) = value {
for (k, v) in more {
o.insert(k, v);
}
}
Value::Object(o)
}
("union", Some(Value::Array(mut a))) => {
if let Value::Array(more) = value {
for v in more {
if !a.contains(&v) {
a.push(v);
}
}
} else if !a.contains(&value) {
a.push(value);
}
Value::Array(a)
}
(_, _) => value,
};
self.vars.insert(key.to_string(), next);
self.touch();
}
pub fn progress(&self) -> Value {
let mut counts: BTreeMap<&str, u32> = BTreeMap::new();
for s in self.steps.values() {
*counts
.entry(match s.status {
StepStatus::Pending => "pending",
StepStatus::Running => "running",
StepStatus::Done => "done",
StepStatus::Failed => "failed",
StepStatus::Skipped => "skipped",
StepStatus::Pruned => "pruned",
StepStatus::Cancelled => "cancelled",
StepStatus::Timeout => "timeout",
StepStatus::Suspended => "suspended",
})
.or_default() += 1;
}
json!(counts)
}
pub fn summary(&self) -> Value {
json!({
"id": self.id, "workflow": self.workflow, "status": self.status, "start": self.start.node,
"steps": self.progress(), "tokens": self.tokens, "created": self.created, "updated": self.updated,
"finished": self.finished, "output": self.output, "error": self.error, "task": self.task, "principal": self.principal,
})
}
}
#[derive(Debug, Clone, PartialEq)]
pub enum Next {
Ready(Vec<String>),
Waiting,
Stalled,
Terminal,
}
pub fn schedule(wf: &Workflow, run: &mut RunState, data: &Data) -> Result<Next, String> {
if run.status.is_terminal() {
return Ok(Next::Terminal);
}
let mut ready = Vec::new();
let mut in_flight = false;
let mut changed = true;
while changed {
changed = false;
for id in wf.topo_order() {
let step = &wf.steps[&id];
let st = run.steps.get(&id).cloned().unwrap_or_default();
match st.status {
StepStatus::Running => {
in_flight = true;
continue;
}
StepStatus::Suspended => {
in_flight = true;
continue;
}
s if s.is_terminal() => continue,
_ => {}
}
if ready.contains(&id) {
continue;
}
if st.forced {
ready.push(id.clone());
continue;
}
if step.depends_on.is_empty()
&& !step.is_start()
&& wf
.steps
.values()
.any(|s| s.field_str("on_timeout") == Some(id.as_str()))
{
continue;
}
let pruned_deps = step
.depends_on
.iter()
.filter(|d| {
run.steps
.get(*d)
.is_some_and(|s| s.status == StepStatus::Pruned)
})
.count();
if !step.depends_on.is_empty() && pruned_deps == step.depends_on.len() {
run.end_step(&id, StepStatus::Pruned, None, None);
changed = true;
continue;
}
let deps_ok = step.depends_on.iter().all(|d| {
run.steps
.get(d)
.is_some_and(|s| s.status.is_satisfied() || s.status == StepStatus::Pruned)
});
let deps_failed = step.depends_on.iter().any(|d| {
run.steps.get(d).is_some_and(|s| {
matches!(
s.status,
StepStatus::Failed | StepStatus::Cancelled | StepStatus::Timeout
)
})
});
if deps_failed {
continue;
}
if !deps_ok {
continue;
}
if let Some(w) = &step.when {
let expr = w.trim().trim_start_matches("CEL:").trim();
let vars: Vec<(&str, &Value)> = data.iter().map(|(k, v)| (k.as_str(), v)).collect();
match crate::cel::eval_bool(expr, &vars) {
Ok(true) => {}
Ok(false) => {
run.end_step(&id, StepStatus::Pruned, None, None);
changed = true;
continue;
}
Err(e) => return Err(format!("step {id:?}: when: {e}")),
}
}
ready.push(id.clone());
}
}
if !ready.is_empty() {
return Ok(Next::Ready(ready));
}
if in_flight {
return Ok(Next::Waiting);
}
Ok(Next::Stalled)
}
pub fn route_failure(
wf: &Workflow,
run: &mut RunState,
step: &Step,
error: &str,
) -> Result<Vec<String>, String> {
match &step.on_error {
OnError::Fail => Err(format!("step {:?} failed: {error}", step.id)),
OnError::Continue => {
let st = run.steps.entry(step.id.clone()).or_default();
st.status = StepStatus::Done;
st.error = Some(error.to_string());
if st.output.is_none() {
st.output = Some(json!({"error": error}));
}
run.touch();
Ok(Vec::new())
}
OnError::Goto(target) => {
if !wf.steps.contains_key(target) {
return Err(format!(
"step {:?}: on_error goto {target:?} does not exist",
step.id
));
}
let st = run.steps.entry(target.clone()).or_default();
st.status = StepStatus::Pending;
st.forced = true;
run.touch();
Ok(vec![target.clone()])
}
}
}
pub fn deadline_passed(run: &RunState) -> bool {
run.deadline_ms.is_some_and(|d| now_ms() >= d)
}
pub fn idempotency_key(run_id: &str, step_id: &str) -> String {
let h = crate::sha::sha256_hex(format!("{run_id}.{step_id}").as_bytes());
h[..32].to_string()
}
pub fn env_view(
instance: &str,
run_id: &str,
instruction: Option<&str>,
prompt: Option<&str>,
) -> Value {
json!({
"instance": instance,
"run": run_id,
"ts": now_ms(),
"instruction": instruction,
"prompt": prompt,
})
}
pub fn render_spec(step: &Step, data: &Data) -> Result<Map<String, Value>, String> {
let mut out = Map::new();
for (k, v) in &step.spec {
if super::model::is_raw_field(&step.kind, k) {
out.insert(k.clone(), v.clone());
continue;
}
out.insert(
k.clone(),
template::render(v, data).map_err(|e| format!("step {:?}: {k}: {e}", step.id))?,
);
}
Ok(out)
}
#[cfg(all(test, feature = "cel"))]
mod tests {
use super::*;
use crate::engine::model::parse_workflow;
#[test]
fn idempotency_keys_are_stable_per_step_and_opaque() {
let a = idempotency_key("run-01ABC", "charge");
assert_eq!(
a,
idempotency_key("run-01ABC", "charge"),
"a retry carries the SAME key"
);
assert_ne!(
a,
idempotency_key("run-01ABC", "refund"),
"another step is another operation"
);
assert_ne!(
a,
idempotency_key("run-02XYZ", "charge"),
"another run is another operation"
);
assert_ne!(
idempotency_key("r", "each[0].call"),
idempotency_key("r", "each[1].call")
);
assert_eq!(a.len(), 32);
assert!(a.chars().all(|c| c.is_ascii_hexdigit()), "hex only: {a}");
assert!(
!a.contains("run-01ABC") && !a.contains("charge"),
"leaks nothing"
);
}
fn start_at(node: &str) -> Start {
Start {
node: node.into(),
payload: json!({}),
ts: 0,
}
}
#[test]
fn an_untaken_branch_prunes_its_tail_but_not_a_live_join() {
let w = parse_workflow(&json!({
"name": "w", "steps": {
"go": {"kind": "once"},
"la": {"kind": "noop", "depends_on": ["go"]},
"ra": {"kind": "noop", "depends_on": ["go"]},
"la2": {"kind": "noop", "depends_on": ["la"]},
"ra2": {"kind": "noop", "depends_on": ["ra"]},
"fin": {"kind": "finish", "depends_on": ["la2", "ra2"], "status": "completed"}
}
}))
.unwrap();
let mut run = RunState::new("r", &w, start_at("go"), json!({}));
run.end_step("ra", StepStatus::Pruned, None, None);
run.end_step("la", StepStatus::Done, None, None);
let data = run.data(env_view("i", "r", None, None), json!({}));
let _ = schedule(&w, &mut run, &data).unwrap();
assert_eq!(
run.steps["ra2"].status,
StepStatus::Pruned,
"the dead branch's tail must be pruned, not run"
);
run.end_step("la2", StepStatus::Done, None, None);
let data = run.data(env_view("i", "r", None, None), json!({}));
match schedule(&w, &mut run, &data).unwrap() {
Next::Ready(r) => assert!(
r.iter().any(|s| s == "fin"),
"a join with one pruned and one live parent must run, got {r:?}"
),
other => panic!("expected fin ready, got {other:?}"),
}
}
#[test]
fn sibling_start_nodes_still_satisfy_their_dependents() {
let w = parse_workflow(&json!({
"name": "w", "steps": {
"a": {"kind": "once"},
"b": {"kind": "manual"},
"work": {"kind": "noop", "depends_on": ["a", "b"]},
"fin": {"kind": "finish", "depends_on": ["work"], "status": "completed"}
}
}))
.unwrap();
let mut run = RunState::new("r", &w, start_at("a"), json!({}));
assert_eq!(run.steps["b"].status, StepStatus::Skipped);
let data = run.data(env_view("i", "r", None, None), json!({}));
match schedule(&w, &mut run, &data).unwrap() {
Next::Ready(r) => assert!(
r.iter().any(|s| s == "work"),
"a step below several start nodes must run when one fired, got {r:?}"
),
other => panic!("expected work ready, got {other:?}"),
}
}
fn wf() -> Workflow {
parse_workflow(&json!({
"name": "w", "steps": {
"s": {"kind": "once"},
"a": {"kind": "noop", "depends_on": ["s"]},
"b": {"kind": "noop", "depends_on": ["s"], "when": "CEL: inputs.go == true"},
"c": {"kind": "noop", "depends_on": ["a", "b"], "on_error": "goto:fix"},
"fix": {"kind": "noop", "depends_on": ["c"]},
"f": {"kind": "finish", "depends_on": ["c"], "status": "completed", "output": "{{vars.x | none}}"}
}
}))
.unwrap()
}
#[cfg(feature = "cel")]
#[test]
fn scheduling_guards_failures_and_terminal_states() {
let w = wf();
let mut run = RunState::new(
"r1",
&w,
Start {
node: "s".into(),
payload: json!({"p": 1}),
ts: 0,
},
json!({"go": false}),
);
assert_eq!(run.steps["s"].status, StepStatus::Done);
assert_eq!(run.steps["s"].output, Some(json!({"p": 1})));
let data = run.data(env_view("i", "r1", None, None), json!({}));
assert_eq!(
schedule(&w, &mut run, &data).unwrap(),
Next::Ready(vec!["a".to_string()])
);
assert_eq!(run.steps["b"].status, StepStatus::Pruned);
run.begin_step("a");
let data = run.data(env_view("i", "r1", None, None), json!({}));
assert_eq!(schedule(&w, &mut run, &data).unwrap(), Next::Waiting);
run.end_step("a", StepStatus::Done, Some(json!("A")), None);
let data = run.data(env_view("i", "r1", None, None), json!({}));
assert_eq!(
schedule(&w, &mut run, &data).unwrap(),
Next::Ready(vec!["c".to_string()])
);
run.begin_step("c");
run.end_step("c", StepStatus::Failed, None, Some("boom".into()));
let next = route_failure(&w, &mut run, w.step("c").unwrap(), "boom").unwrap();
assert_eq!(next, vec!["fix".to_string()]);
run.begin_step("fix");
run.end_step("fix", StepStatus::Done, None, None);
let data = run.data(env_view("i", "r1", None, None), json!({}));
assert_eq!(schedule(&w, &mut run, &data).unwrap(), Next::Stalled);
run.finish(RunStatus::Stalled, None, Some("stalled".into()));
assert!(run.status.is_terminal());
let data = run.data(env_view("i", "r1", None, None), json!({}));
assert_eq!(schedule(&w, &mut run, &data).unwrap(), Next::Terminal);
let mut run2 = RunState::new(
"r2",
&w,
Start {
node: "s".into(),
payload: json!({}),
ts: 0,
},
json!({"go": true}),
);
let mut c = w.step("c").unwrap().clone();
c.on_error = OnError::Continue;
run2.begin_step("c");
run2.end_step("c", StepStatus::Failed, None, Some("e".into()));
assert!(route_failure(&w, &mut run2, &c, "e").unwrap().is_empty());
assert_eq!(run2.steps["c"].status, StepStatus::Done);
assert_eq!(run2.steps["c"].error.as_deref(), Some("e"));
let mut a = w.step("a").unwrap().clone();
a.on_error = OnError::Fail;
assert!(route_failure(&w, &mut run2, &a, "e").is_err());
}
#[test]
fn vars_reducers_and_serialization() {
let w = wf();
let mut run = RunState::new("r", &w, Start::default(), json!({}));
run.write_var("l", json!([1]), "overwrite");
run.write_var("l", json!(2), "append");
run.write_var("l", json!([3, 4]), "append");
assert_eq!(run.vars["l"], json!([1, 2, 3, 4]));
run.write_var("l", json!([4, 5]), "union");
assert_eq!(run.vars["l"], json!([1, 2, 3, 4, 5]));
run.write_var("o", json!({"a": 1}), "overwrite");
run.write_var("o", json!({"b": 2}), "merge");
assert_eq!(run.vars["o"], json!({"a": 1, "b": 2}));
run.write_var("o", json!(7), "overwrite");
assert_eq!(run.vars["o"], json!(7));
let v = serde_json::to_value(&run).unwrap();
let back: RunState = serde_json::from_value(v).unwrap();
assert_eq!(back.vars, run.vars);
assert!(!back.dirty);
assert_eq!(back.summary()["workflow"], json!("w"));
let data = run.data(
env_view("inst", "r", Some("brief"), None),
json!({"k": "v"}),
);
let mut s = w.step("f").unwrap().clone();
s.spec
.insert("extra".into(), json!("{{env.instruction}}/{{memory.k}}"));
let rendered = render_spec(&s, &data).unwrap();
assert_eq!(rendered["output"], json!("none"));
assert_eq!(rendered["extra"], json!("brief/v"));
assert!(!deadline_passed(&run));
}
}