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