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