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