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