ironflow_engine/context/accessors.rs
1//! Constructors and plain accessors for [`WorkflowContext`].
2//!
3//! Everything here reads or writes a single field: the run identity, the
4//! wiring the engine attaches before execution, the run-level budget, and the
5//! counters a handler can inspect while it runs.
6
7use std::collections::HashMap;
8use std::sync::Arc;
9
10use rust_decimal::Decimal;
11use serde_json::{Value, from_value};
12use tracing::error;
13use uuid::Uuid;
14
15use ironflow_core::decision::DecisionProvider;
16use ironflow_core::provider::AgentProvider;
17use ironflow_core::trace_context::WorkflowTraceContext;
18use ironflow_store::error::StoreError;
19use ironflow_store::models::Step;
20use ironflow_store::store::Store;
21#[cfg(feature = "secret-store")]
22use ironflow_store::workflow_secrets::ScopedSecretStore;
23
24use crate::artifact::ArtifactSink;
25use crate::error::EngineError;
26use crate::executor::{StepInterceptor, StepResult};
27use crate::guard::{SharedGuardState, WorkflowGuardConfig};
28use crate::log_sender::LogSender;
29use crate::notify::WorkflowEventBus;
30#[cfg(not(feature = "secret-store"))]
31use crate::operation::NoopSecretResolver;
32use crate::operation::{OperationContext, SecretResolver};
33use crate::plan::{SharedPlanRecorder, lock_plan};
34
35use super::{HandlerResolver, WorkflowContext};
36
37impl WorkflowContext {
38 /// Create a new context for a run.
39 ///
40 /// Not typically called directly — the [`Engine`](crate::engine::Engine)
41 /// creates this when executing a [`WorkflowHandler`](crate::handler::WorkflowHandler).
42 pub fn new(
43 run_id: Uuid,
44 workflow_name: String,
45 store: Arc<dyn Store>,
46 provider: Arc<dyn AgentProvider>,
47 ) -> Self {
48 let trace_context = WorkflowTraceContext::from_workflow_run_id(&run_id.to_string());
49 Self {
50 run_id,
51 root_run_id: run_id,
52 workflow_name,
53 store,
54 provider,
55 decision_provider: None,
56 handler_resolver: None,
57 position: 0,
58 last_step_ids: Vec::new(),
59 total_cost_usd: Decimal::ZERO,
60 total_duration_ms: 0,
61 max_cost_usd: None,
62 inherited_cost_usd: Decimal::ZERO,
63 replay_steps: HashMap::new(),
64 replay_wave_steps: HashMap::new(),
65 granted_approvals: HashMap::new(),
66 answered_inputs: HashMap::new(),
67 attempt: 1,
68 carried_duration_ms: 0,
69 log_sender: None,
70 artifact_sink: None,
71 has_allowed_failure: false,
72 error_handlers: Vec::new(),
73 guard_state: None,
74 guard_config: None,
75 step_results: Vec::new(),
76 event_bus: None,
77 interceptor: None,
78 trace_context,
79 operation_ctx: None,
80 plan: None,
81 }
82 }
83
84 /// Create a new context with a handler resolver for sub-workflow support.
85 ///
86 /// The resolver is called when [`workflow`](Self::workflow) is invoked to
87 /// look up registered handlers by name.
88 pub(crate) fn with_handler_resolver(
89 run_id: Uuid,
90 workflow_name: String,
91 store: Arc<dyn Store>,
92 provider: Arc<dyn AgentProvider>,
93 resolver: HandlerResolver,
94 ) -> Self {
95 let trace_context = WorkflowTraceContext::from_workflow_run_id(&run_id.to_string());
96 Self {
97 run_id,
98 root_run_id: run_id,
99 workflow_name,
100 store,
101 provider,
102 decision_provider: None,
103 handler_resolver: Some(resolver),
104 position: 0,
105 last_step_ids: Vec::new(),
106 total_cost_usd: Decimal::ZERO,
107 total_duration_ms: 0,
108 max_cost_usd: None,
109 inherited_cost_usd: Decimal::ZERO,
110 replay_steps: HashMap::new(),
111 replay_wave_steps: HashMap::new(),
112 granted_approvals: HashMap::new(),
113 answered_inputs: HashMap::new(),
114 attempt: 1,
115 carried_duration_ms: 0,
116 log_sender: None,
117 artifact_sink: None,
118 has_allowed_failure: false,
119 error_handlers: Vec::new(),
120 guard_state: None,
121 guard_config: None,
122 step_results: Vec::new(),
123 event_bus: None,
124 interceptor: None,
125 trace_context,
126 operation_ctx: None,
127 plan: None,
128 }
129 }
130
131 /// Attach a log sender for real-time step output streaming.
132 pub fn set_log_sender(&mut self, sender: LogSender) {
133 self.log_sender = Some(sender);
134 }
135
136 /// Attach the backend that stores and serves artifact bytes.
137 ///
138 /// Without one, any step that declares an output or calls
139 /// [`put_artifact`](Self::put_artifact) fails with
140 /// [`EngineError::ArtifactsUnavailable`]. Every other step is unaffected,
141 /// so an existing deployment keeps working until artifacts are configured.
142 ///
143 /// # Examples
144 ///
145 /// ```no_run
146 /// use std::sync::Arc;
147 ///
148 /// use ironflow_engine::artifact::ArtifactSink;
149 /// use ironflow_engine::context::WorkflowContext;
150 ///
151 /// # fn example(ctx: &mut WorkflowContext, sink: Arc<dyn ArtifactSink>) {
152 /// ctx.set_artifact_sink(sink);
153 /// # }
154 /// ```
155 pub fn set_artifact_sink(&mut self, sink: Arc<dyn ArtifactSink>) {
156 self.artifact_sink = Some(sink);
157 }
158
159 /// Return the W3C trace context for this workflow run.
160 ///
161 /// The trace context is derived from the run ID and can be used to
162 /// correlate spans across distributed services. Each step automatically
163 /// receives a [`child`](WorkflowTraceContext::child) context.
164 pub fn trace_context(&self) -> &WorkflowTraceContext {
165 &self.trace_context
166 }
167
168 /// Attach a workflow guard configuration and shared state.
169 ///
170 /// When set, the guard is checked before every sub-workflow invocation.
171 /// The shared state is propagated to child workflows so that limits
172 /// apply globally across the entire run tree.
173 ///
174 /// # Examples
175 ///
176 /// ```no_run
177 /// use ironflow_engine::context::WorkflowContext;
178 /// use ironflow_engine::guard::{WorkflowGuardConfig, new_shared_guard_state};
179 ///
180 /// # fn example(ctx: &mut WorkflowContext) {
181 /// ctx.set_guard(WorkflowGuardConfig::default(), new_shared_guard_state());
182 /// # }
183 /// ```
184 pub fn set_guard(&mut self, config: WorkflowGuardConfig, state: SharedGuardState) {
185 self.guard_config = Some(config);
186 self.guard_state = Some(state);
187 }
188
189 /// The current guard configuration, if any.
190 pub fn guard_config(&self) -> Option<&WorkflowGuardConfig> {
191 self.guard_config.as_ref()
192 }
193
194 /// Attach a [`WorkflowEventBus`] for per-run real-time monitoring.
195 ///
196 /// When set, step transitions automatically publish
197 /// [`WorkflowEvent`](crate::notify::WorkflowEvent)s to the bus.
198 pub fn set_event_bus(&mut self, bus: WorkflowEventBus) {
199 self.event_bus = Some(bus);
200 }
201
202 /// Attach a [`StepInterceptor`] that resolves steps without executing them.
203 ///
204 /// Wired by the [`Engine`](crate::engine::Engine) from
205 /// [`Engine::with_step_interceptor`](crate::engine::Engine::with_step_interceptor).
206 /// Intended for tests: see [`crate::testing::TestEngine`].
207 ///
208 /// # Examples
209 ///
210 /// ```no_run
211 /// use std::sync::Arc;
212 ///
213 /// use ironflow_engine::context::WorkflowContext;
214 /// use ironflow_engine::executor::StepInterceptor;
215 ///
216 /// # fn example(ctx: &mut WorkflowContext, interceptor: Arc<dyn StepInterceptor>) {
217 /// ctx.set_step_interceptor(interceptor);
218 /// # }
219 /// ```
220 pub fn set_step_interceptor(&mut self, interceptor: Arc<dyn StepInterceptor>) {
221 self.interceptor = Some(interceptor);
222 }
223
224 /// The step interceptor attached to this context, if any.
225 pub(crate) fn step_interceptor(&self) -> Option<&Arc<dyn StepInterceptor>> {
226 self.interceptor.as_ref()
227 }
228
229 /// Attach a [`DecisionProvider`] backend for `ctx.decision(...)` steps.
230 ///
231 /// Not typically called directly -- the [`Engine`](crate::engine::Engine)
232 /// wires this from [`Engine::with_decision_provider`](crate::engine::Engine::with_decision_provider).
233 pub fn set_decision_provider(&mut self, provider: Arc<dyn DecisionProvider>) {
234 self.decision_provider = Some(provider);
235 }
236
237 /// Attach a plan recorder, switching this context to plan mode.
238 ///
239 /// Every step method then records its intent instead of executing it.
240 pub(crate) fn set_plan(&mut self, plan: SharedPlanRecorder) {
241 self.plan = Some(plan);
242 }
243
244 /// The shared recorder, when this context is planning.
245 pub(crate) fn plan(&self) -> Option<&SharedPlanRecorder> {
246 self.plan.as_ref()
247 }
248
249 /// Whether this context records a plan instead of executing steps.
250 ///
251 /// # Examples
252 ///
253 /// ```no_run
254 /// use ironflow_engine::context::WorkflowContext;
255 ///
256 /// # fn example(ctx: &WorkflowContext) {
257 /// if ctx.is_planning() {
258 /// // No command runs, no request is sent: only the plan is recorded.
259 /// }
260 /// # }
261 /// ```
262 pub fn is_planning(&self) -> bool {
263 self.plan.is_some()
264 }
265
266 /// Seed the context with the run's attempt number and the totals already
267 /// accumulated by previous attempts.
268 ///
269 /// Called by the engine before executing a handler. Steps created by this
270 /// context belong to `attempt`, and the cost and duration it reports at the
271 /// end cover the whole run, not just this attempt.
272 pub(crate) fn carry_over_run_totals(
273 &mut self,
274 attempt: u32,
275 cost_usd: Decimal,
276 duration_ms: u64,
277 ) {
278 self.attempt = attempt;
279 self.total_cost_usd = cost_usd;
280 self.carried_duration_ms = duration_ms;
281 }
282
283 /// Wall-clock duration already recorded on the run by previous attempts.
284 pub(crate) fn carried_duration_ms(&self) -> u64 {
285 self.carried_duration_ms
286 }
287
288 /// The run attempt this context is executing (1-based).
289 pub fn attempt(&self) -> u32 {
290 self.attempt
291 }
292
293 /// Set the cumulative cost cap enforced before every agent step.
294 ///
295 /// Called by the [`Engine`](crate::engine::Engine) with the run's persisted
296 /// `max_cost_usd`. `None` disables the check.
297 ///
298 /// # Examples
299 ///
300 /// ```no_run
301 /// use ironflow_engine::context::WorkflowContext;
302 /// use rust_decimal::Decimal;
303 ///
304 /// # fn example(ctx: &mut WorkflowContext) {
305 /// ctx.set_max_cost_usd(Some(Decimal::new(200, 2))); // $2.00
306 /// # }
307 /// ```
308 pub fn set_max_cost_usd(&mut self, cap: Option<Decimal>) {
309 self.max_cost_usd = cap;
310 }
311
312 /// The cumulative cost cap of this run, if any.
313 pub fn max_cost_usd(&self) -> Option<Decimal> {
314 self.max_cost_usd
315 }
316
317 /// Total cost charged against the cap: this run plus every ancestor run.
318 ///
319 /// For a top-level run this equals [`total_cost_usd`](Self::total_cost_usd).
320 /// For a sub-workflow it also includes what the parent chain already spent.
321 pub fn charged_cost_usd(&self) -> Decimal {
322 self.inherited_cost_usd + self.total_cost_usd
323 }
324
325 /// Reject the upcoming agent work when it would cross the run's cost cap.
326 ///
327 /// `step_budget` is the declared budget of the step (or the sum of budgets
328 /// for a parallel wave). Called *before* any step record is created so a
329 /// refused run never launches the work it could not afford.
330 ///
331 /// # Errors
332 ///
333 /// Returns [`EngineError::RunBudgetExceeded`] when
334 /// `charged_cost + step_budget` exceeds the cap.
335 pub(super) fn check_run_budget(&self, step_budget: Decimal) -> Result<(), EngineError> {
336 let Some(limit) = self.max_cost_usd else {
337 return Ok(());
338 };
339
340 let spent = self.charged_cost_usd();
341 if spent + step_budget <= limit {
342 return Ok(());
343 }
344
345 error!(
346 run_id = %self.run_id,
347 limit_usd = %limit,
348 spent_usd = %spent,
349 step_budget_usd = %step_budget,
350 "run cost cap reached, refusing agent step"
351 );
352
353 Err(EngineError::RunBudgetExceeded {
354 run_id: self.run_id,
355 limit_usd: limit,
356 spent_usd: spent,
357 step_budget_usd: step_budget,
358 })
359 }
360
361 /// The run ID this context is executing for.
362 pub fn run_id(&self) -> Uuid {
363 self.run_id
364 }
365
366 /// The top-level run this context belongs to: [`run_id`](Self::run_id)
367 /// itself, or, inside a sub-workflow, the run that started the chain.
368 ///
369 /// The engine stamps it on agent pods as
370 /// [`LABEL_ROOT_RUN_ID`](ironflow_core::provider::LABEL_ROOT_RUN_ID); set
371 /// the same label on a pod you create yourself (`PodRun`, `JobRun`) so
372 /// that a retry of the top-level run deletes what a dead attempt left.
373 ///
374 /// # Examples
375 ///
376 /// ```no_run
377 /// use ironflow_core::provider::LABEL_ROOT_RUN_ID;
378 /// use ironflow_engine::context::WorkflowContext;
379 ///
380 /// fn root_label(ctx: &WorkflowContext) -> (&'static str, String) {
381 /// (LABEL_ROOT_RUN_ID, ctx.root_run_id().to_string())
382 /// }
383 /// ```
384 pub fn root_run_id(&self) -> Uuid {
385 self.root_run_id
386 }
387
388 /// The workflow name this run belongs to.
389 pub fn workflow_name(&self) -> &str {
390 &self.workflow_name
391 }
392
393 /// Accumulated cost across all executed steps so far.
394 pub fn total_cost_usd(&self) -> Decimal {
395 self.total_cost_usd
396 }
397
398 /// Whether at least one `allow_failure` step failed during this run.
399 pub fn has_allowed_failure(&self) -> bool {
400 self.has_allowed_failure
401 }
402
403 /// Accumulated duration across all executed steps so far.
404 pub fn total_duration_ms(&self) -> u64 {
405 self.total_duration_ms
406 }
407
408 /// Enriched results of all completed steps in execution order.
409 pub fn step_results(&self) -> &[StepResult] {
410 &self.step_results
411 }
412
413 /// Return a [`ScopedSecretStore`] scoped to this workflow.
414 ///
415 /// Secrets are namespaced under `workflows/<uuid>/` where the UUID is
416 /// deterministically derived from the workflow name (UUID v5). This means
417 /// all runs of the same workflow share the same secret namespace.
418 ///
419 /// Requires the `secret-store` feature.
420 ///
421 /// # Examples
422 ///
423 /// ```no_run
424 /// use ironflow_engine::context::WorkflowContext;
425 /// use ironflow_engine::error::EngineError;
426 ///
427 /// # async fn example(ctx: &WorkflowContext) -> Result<(), EngineError> {
428 /// let secrets = ctx.secrets();
429 /// secrets.set("api_token", "sk-ant-12345").await.map_err(EngineError::Store)?;
430 /// let token = secrets.get("api_token").await.map_err(EngineError::Store)?;
431 /// # Ok(())
432 /// # }
433 /// ```
434 #[cfg(feature = "secret-store")]
435 pub fn secrets(&self) -> ScopedSecretStore {
436 let workflow_uuid = Uuid::new_v5(&Uuid::NAMESPACE_OID, self.workflow_name.as_bytes());
437 ScopedSecretStore::for_workflow(workflow_uuid, self.store.clone())
438 }
439
440 pub(super) fn ensure_operation_ctx(&mut self) -> &OperationContext {
441 self.operation_ctx.get_or_insert_with(|| {
442 #[cfg(feature = "secret-store")]
443 let secrets: Arc<dyn SecretResolver> = {
444 let workflow_uuid =
445 Uuid::new_v5(&Uuid::NAMESPACE_OID, self.workflow_name.as_bytes());
446 Arc::new(ScopedSecretStore::for_workflow(
447 workflow_uuid,
448 self.store.clone(),
449 ))
450 };
451 #[cfg(not(feature = "secret-store"))]
452 let secrets: Arc<dyn SecretResolver> = Arc::new(NoopSecretResolver);
453
454 OperationContext::new(secrets)
455 })
456 }
457
458 /// Access the store directly (advanced usage).
459 pub fn store(&self) -> &Arc<dyn Store> {
460 &self.store
461 }
462
463 /// Get and increment the current position counter.
464 pub(crate) fn next_position(&mut self) -> u32 {
465 let pos = self.position;
466 self.position += 1;
467 pos
468 }
469
470 /// Access the replay steps from a previous execution.
471 pub(crate) fn replay_steps(&self) -> &HashMap<u32, Step> {
472 &self.replay_steps
473 }
474
475 /// Set the last step IDs (for dependency tracking).
476 pub(crate) fn set_last_step_ids(&mut self, ids: Vec<Uuid>) {
477 self.last_step_ids = ids;
478 }
479
480 /// Access the payload that triggered this run.
481 ///
482 /// Fetches the run from the store and returns its payload.
483 ///
484 /// # Errors
485 ///
486 /// Returns [`EngineError::Store`] if the run is not found.
487 pub async fn payload(&self) -> Result<Value, EngineError> {
488 // In plan mode the payload comes from the recorder: planning must not
489 // create, or even read, a run record.
490 if let Some(plan) = &self.plan {
491 return Ok(lock_plan(plan).payload());
492 }
493
494 let run = self
495 .store
496 .get_run(self.run_id)
497 .await?
498 .ok_or(EngineError::Store(StoreError::RunNotFound(self.run_id)))?;
499 Ok(run.payload)
500 }
501
502 /// Deserialize the run payload into a typed input struct.
503 ///
504 /// Shorthand for `serde_json::from_value(ctx.payload().await?)`.
505 ///
506 /// # Errors
507 ///
508 /// Returns [`EngineError::Store`] if the run is not found, or
509 /// [`EngineError::Serialization`] if the payload does not match `T`.
510 ///
511 /// # Examples
512 ///
513 /// ```no_run
514 /// # use ironflow_engine::context::WorkflowContext;
515 /// # use ironflow_engine::error::EngineError;
516 /// use serde::Deserialize;
517 ///
518 /// #[derive(Deserialize)]
519 /// struct DeployInput {
520 /// environment: String,
521 /// dry_run: Option<bool>,
522 /// }
523 ///
524 /// # async fn example(ctx: &WorkflowContext) -> Result<(), EngineError> {
525 /// let input: DeployInput = ctx.input().await?;
526 /// # Ok(())
527 /// # }
528 /// ```
529 pub async fn input<T: serde::de::DeserializeOwned>(&self) -> Result<T, EngineError> {
530 let payload = self.payload().await?;
531 from_value(payload).map_err(EngineError::Serialization)
532 }
533}