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            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}