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