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