1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
//! Typed machine-decision step for [`WorkflowContext`].
//!
//! Holds the public [`decision`](WorkflowContext::decision) entry point and the
//! replay, execute and escalate paths it dispatches to. As a descendant module
//! of `context`, it can access `WorkflowContext`'s private fields.
//!
//! The decision step is a hybrid of an agent step (it calls a provider, costs
//! money, and stores typed output) and an approval gate (a low-confidence answer
//! suspends the run and replays its stored answers on resume).
use std::collections::BTreeMap;
use chrono::Utc;
use serde_json::{Value, from_value, to_value};
use tracing::info;
use uuid::Uuid;
use ironflow_core::decision::{DecisionOutput, DecisionUsage};
use ironflow_store::models::{NewStep, StepKind, StepStatus, StepUpdate, step_trace_id};
use crate::config::DecisionConfig;
use crate::context::WorkflowContext;
use crate::context::lifecycle::check_replay_identity;
use crate::decision::DecisionAnswers;
use crate::error::EngineError;
use crate::executor::{DecisionExecution, StepArtifacts, StepOutput, StepResult, execute_decision};
use crate::notify::{
WorkflowApprovalRequiredEvent, WorkflowEvent, WorkflowStepCompletedEvent,
WorkflowStepStartedEvent,
};
use crate::plan::lock_plan;
impl WorkflowContext {
/// Execute a typed machine-decision step (System One / Jev).
///
/// The questions come from `T`, set with
/// [`DecisionConfig::answers`], and the answers are returned as a `T`. See
/// [`crate::decision`]. When `escalate_below` is set and any answer falls
/// below it, the run suspends with [`EngineError::ApprovalRequired`] and
/// replays the stored answers on resume without re-calling the provider.
///
/// While planning no provider is called and no answer exists, so reading
/// them fails with [`EngineError::Decision`] and the plan stops at this
/// step.
///
/// # Errors
///
/// [`EngineError::NoDecisionProvider`], [`EngineError::ApprovalRequired`],
/// [`EngineError::Operation`], [`EngineError::Decision`] when an answer
/// does not fit `T`, or [`EngineError::ReplayDivergence`] when the step
/// recorded at this position has a different name or kind.
///
/// # Examples
///
/// ```no_run
/// use ironflow_engine::config::DecisionConfig;
/// use ironflow_engine::context::WorkflowContext;
/// use ironflow_engine::decision::DecisionAnswers;
/// use ironflow_engine::error::EngineError;
///
/// #[derive(DecisionAnswers)]
/// struct Urgency {
/// #[noul("Does this convey urgency?")]
/// urgent: f64,
/// }
///
/// # async fn example(ctx: &mut WorkflowContext) -> Result<(), EngineError> {
/// let urgency = ctx
/// .decision("urgency", DecisionConfig::new("Payouts fail since 3 days").answers::<Urgency>())
/// .await?;
/// if urgency.urgent > 0.8 {
/// // page someone
/// }
/// # Ok(())
/// # }
/// ```
///
/// A config without questions is not a decision:
///
/// ```compile_fail,E0277
/// # use ironflow_engine::config::DecisionConfig;
/// # use ironflow_engine::context::WorkflowContext;
/// # use ironflow_engine::error::EngineError;
/// # async fn example(ctx: &mut WorkflowContext) -> Result<(), EngineError> {
/// ctx.decision("urgency", DecisionConfig::new("state")).await?;
/// # Ok(())
/// # }
/// ```
pub async fn decision<T: DecisionAnswers>(
&mut self,
name: &str,
config: DecisionConfig<T>,
) -> Result<T, EngineError> {
let output = self.decision_output(name, config.erase()).await?;
Ok(T::from_output(&output)?)
}
/// Run, replay or plan a decision step and return the raw answers.
async fn decision_output(
&mut self,
name: &str,
config: DecisionConfig,
) -> Result<DecisionOutput, EngineError> {
// Plan mode: record the step and return an empty answer set. No
// provider is called, so the step costs nothing while planning.
if let Some(plan) = self.plan().cloned() {
self.position += 1;
let mut recorder = lock_plan(&plan);
if recorder.record(name, StepKind::Decision, &self.workflow_name, None) {
recorder.set_last(vec![name.to_string()]);
}
return Ok(DecisionOutput {
model: None,
answers: BTreeMap::new(),
usage: DecisionUsage::default(),
});
}
if let Some(output) = self.decision_replay(name, &config).await? {
return Ok(output);
}
self.decision_execute(name, config).await
}
/// Replay a decision step recorded in this attempt, if any.
///
/// Returns `Ok(Some(output))` when the step at the current position was
/// already decided (its answers are returned as-is, without re-calling the
/// provider), advancing the position. Returns `Ok(None)` when there is
/// nothing to replay, leaving the position untouched for a fresh execution.
async fn decision_replay(
&mut self,
name: &str,
_config: &DecisionConfig,
) -> Result<Option<DecisionOutput>, EngineError> {
let position = self.position;
let Some(existing) = self.replay_steps.get(&position).cloned() else {
return Ok(None);
};
check_replay_identity(&existing, position, name, &StepKind::Decision)?;
self.position += 1;
let stored: DecisionOutput = existing
.output
.clone()
.ok_or_else(|| {
EngineError::StepConfig(format!(
"decision step '{name}' has no stored output to replay"
))
})
.and_then(|v| from_value(v).map_err(EngineError::from))?;
// An escalated decision suspended in `AwaitingApproval`; the handler only
// re-runs on an approved resume, so mark it completed and continue.
if existing.status.state == StepStatus::AwaitingApproval {
self.store
.update_step(
existing.id,
StepUpdate {
status: Some(StepStatus::Completed),
completed_at: Some(Utc::now()),
..StepUpdate::default()
},
)
.await?;
info!(
run_id = %self.run_id,
step = %name,
position,
"decision step replayed (approved after escalation)"
);
} else {
info!(
run_id = %self.run_id,
step = %name,
position,
"decision step replayed from previous execution"
);
}
// Do not re-add cost/duration here. A replay only happens on a resume
// within the same attempt, where `carry_over_run_totals` has already
// seeded `total_cost_usd`/`total_duration_ms` from the run totals the
// suspend snapshot persisted -- totals that already include this step.
// Adding them again would double-count the escalated decision.
self.last_step_ids = vec![existing.id];
Ok(Some(stored))
}
/// Execute a fresh decision step: call the provider, persist the answers, and
/// either complete or escalate to a human approval gate.
async fn decision_execute(
&mut self,
name: &str,
config: DecisionConfig,
) -> Result<DecisionOutput, EngineError> {
self.check_guard_timeout()?;
let position = self.position;
self.position += 1;
let provider =
self.decision_provider
.clone()
.ok_or_else(|| EngineError::NoDecisionProvider {
step: name.to_string(),
})?;
let trace_id = step_trace_id(self.run_id, name, position);
let step = self
.store
.create_step(NewStep {
run_id: self.run_id,
trace_id,
name: name.to_string(),
kind: StepKind::Decision,
position,
input: Some(to_value(&config)?),
is_error_handler: false,
})
.await?;
self.start_step(step.id, Utc::now()).await?;
if let Some(ref bus) = self.event_bus {
bus.publish(
self.run_id,
WorkflowEvent::StepStarted(WorkflowStepStartedEvent {
step_name: name.to_string(),
step_index: position,
timestamp: Utc::now(),
}),
);
}
let execution = match execute_decision(&provider, &config).await {
Ok(execution) => execution,
Err(err) => {
self.fail_step(step.id, &err).await;
return Err(err);
}
};
// Impute cost and duration to the run, like an agent step.
self.total_cost_usd += execution.cost_usd;
self.total_duration_ms += execution.duration_ms;
let output_value = to_value(&execution.output)?;
let escalated = config
.escalate_below
.zip(execution.output.min_confidence())
.map(|(threshold, min)| min < threshold)
.unwrap_or(false);
if escalated {
return self
.decision_escalate(name, position, step.id, &config, &execution, output_value)
.await;
}
let step_output = StepOutput {
output: output_value.clone(),
duration_ms: execution.duration_ms,
cost_usd: execution.cost_usd,
input_tokens: Some(execution.input_tokens),
cache_read_input_tokens: None,
cache_creation_input_tokens: None,
output_tokens: Some(execution.output_tokens),
model: execution.output.model.as_ref().map(ToString::to_string),
debug_messages: None,
artifacts: StepArtifacts::default(),
};
let completed_at = Utc::now();
self.store
.update_step(
step.id,
StepUpdate {
status: Some(StepStatus::Completed),
output: Some(output_value),
duration_ms: Some(execution.duration_ms),
cost_usd: Some(execution.cost_usd),
input_tokens: Some(execution.input_tokens),
output_tokens: Some(execution.output_tokens),
completed_at: Some(completed_at),
..StepUpdate::default()
},
)
.await?;
self.step_results
.push(StepResult::from_success(trace_id, name, &step_output));
self.persist_progress().await;
self.last_step_ids = vec![step.id];
info!(
run_id = %self.run_id,
step = %name,
trace_id = %trace_id,
cost_usd = %execution.cost_usd,
"decision step completed"
);
if let Some(ref bus) = self.event_bus {
bus.publish(
self.run_id,
WorkflowEvent::StepCompleted(WorkflowStepCompletedEvent {
step_name: name.to_string(),
step_index: position,
duration_ms: execution.duration_ms,
output_summary: None,
}),
);
}
Ok(execution.output)
}
/// Persist an escalated decision, suspend the run, and return
/// [`EngineError::ApprovalRequired`].
async fn decision_escalate(
&mut self,
name: &str,
position: u32,
step_id: Uuid,
config: &DecisionConfig,
execution: &DecisionExecution,
output_value: Value,
) -> Result<DecisionOutput, EngineError> {
let threshold = config.escalate_below.unwrap_or_default();
let min = execution.output.min_confidence().unwrap_or_default();
self.store
.update_step(
step_id,
StepUpdate {
status: Some(StepStatus::AwaitingApproval),
output: Some(output_value),
duration_ms: Some(execution.duration_ms),
cost_usd: Some(execution.cost_usd),
input_tokens: Some(execution.input_tokens),
output_tokens: Some(execution.output_tokens),
..StepUpdate::default()
},
)
.await?;
self.last_step_ids = vec![step_id];
info!(
run_id = %self.run_id,
step = %name,
position,
confidence = min,
threshold,
"decision escalated to human approval"
);
if let Some(ref bus) = self.event_bus {
bus.publish(
self.run_id,
WorkflowEvent::ApprovalRequired(WorkflowApprovalRequiredEvent {
step_name: name.to_string(),
step_index: position,
approval_id: step_id,
}),
);
}
Err(EngineError::ApprovalRequired {
run_id: self.run_id,
step_id,
message: format!(
"decision '{name}' escalated: confidence {min:.3} below threshold {threshold:.3}"
),
})
}
}