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