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