Skip to main content

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