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