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