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