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