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