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