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