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