pub struct WorkflowContext { /* private fields */ }Expand description
Execution context for a single workflow run.
Tracks the current step position and provides convenience methods for executing operations with automatic persistence.
§Examples
use ironflow_engine::context::WorkflowContext;
use ironflow_engine::config::ShellConfig;
use ironflow_engine::error::EngineError;
let result = ctx.shell("greet", ShellConfig::new("echo hello")).await?;
assert!(result.stdout().contains("hello"));Implementations§
Source§impl WorkflowContext
impl WorkflowContext
Sourcepub fn new(
run_id: Uuid,
workflow_name: String,
store: Arc<dyn Store>,
provider: Arc<dyn AgentProvider>,
) -> Self
pub fn new( run_id: Uuid, workflow_name: String, store: Arc<dyn Store>, provider: Arc<dyn AgentProvider>, ) -> Self
Create a new context for a run.
Not typically called directly — the Engine
creates this when executing a WorkflowHandler.
Sourcepub fn set_log_sender(&mut self, sender: LogSender)
pub fn set_log_sender(&mut self, sender: LogSender)
Attach a log sender for real-time step output streaming.
Sourcepub fn set_artifact_sink(&mut self, sink: Arc<dyn ArtifactSink>)
pub fn set_artifact_sink(&mut self, sink: Arc<dyn ArtifactSink>)
Attach the backend that stores and serves artifact bytes.
Without one, any step that declares an output or calls
put_artifact fails with
EngineError::ArtifactsUnavailable. Every other step is unaffected,
so an existing deployment keeps working until artifacts are configured.
§Examples
use std::sync::Arc;
use ironflow_engine::artifact::ArtifactSink;
use ironflow_engine::context::WorkflowContext;
ctx.set_artifact_sink(sink);Sourcepub fn trace_context(&self) -> &WorkflowTraceContext
pub fn trace_context(&self) -> &WorkflowTraceContext
Return the W3C trace context for this workflow run.
The trace context is derived from the run ID and can be used to
correlate spans across distributed services. Each step automatically
receives a child context.
Sourcepub fn set_guard(
&mut self,
config: WorkflowGuardConfig,
state: SharedGuardState,
)
pub fn set_guard( &mut self, config: WorkflowGuardConfig, state: SharedGuardState, )
Attach a workflow guard configuration and shared state.
When set, the guard is checked before every sub-workflow invocation. The shared state is propagated to child workflows so that limits apply globally across the entire run tree.
§Examples
use ironflow_engine::context::WorkflowContext;
use ironflow_engine::guard::{WorkflowGuardConfig, new_shared_guard_state};
ctx.set_guard(WorkflowGuardConfig::default(), new_shared_guard_state());Sourcepub fn guard_config(&self) -> Option<&WorkflowGuardConfig>
pub fn guard_config(&self) -> Option<&WorkflowGuardConfig>
The current guard configuration, if any.
Sourcepub fn set_event_bus(&mut self, bus: WorkflowEventBus)
pub fn set_event_bus(&mut self, bus: WorkflowEventBus)
Attach a WorkflowEventBus for per-run real-time monitoring.
When set, step transitions automatically publish
WorkflowEvents to the bus.
Sourcepub fn set_step_interceptor(&mut self, interceptor: Arc<dyn StepInterceptor>)
pub fn set_step_interceptor(&mut self, interceptor: Arc<dyn StepInterceptor>)
Attach a StepInterceptor that resolves steps without executing them.
Wired by the Engine from
Engine::with_step_interceptor.
Intended for tests: see crate::testing::TestEngine.
§Examples
use std::sync::Arc;
use ironflow_engine::context::WorkflowContext;
use ironflow_engine::executor::StepInterceptor;
ctx.set_step_interceptor(interceptor);Sourcepub fn set_decision_provider(&mut self, provider: Arc<dyn DecisionProvider>)
pub fn set_decision_provider(&mut self, provider: Arc<dyn DecisionProvider>)
Attach a DecisionProvider backend for ctx.decision(...) steps.
Not typically called directly – the Engine
wires this from Engine::with_decision_provider.
Sourcepub fn is_planning(&self) -> bool
pub fn is_planning(&self) -> bool
Whether this context records a plan instead of executing steps.
§Examples
use ironflow_engine::context::WorkflowContext;
if ctx.is_planning() {
// No command runs, no request is sent: only the plan is recorded.
}Sourcepub fn set_max_cost_usd(&mut self, cap: Option<Decimal>)
pub fn set_max_cost_usd(&mut self, cap: Option<Decimal>)
Sourcepub fn max_cost_usd(&self) -> Option<Decimal>
pub fn max_cost_usd(&self) -> Option<Decimal>
The cumulative cost cap of this run, if any.
Sourcepub fn charged_cost_usd(&self) -> Decimal
pub fn charged_cost_usd(&self) -> Decimal
Total cost charged against the cap: this run plus every ancestor run.
For a top-level run this equals total_cost_usd.
For a sub-workflow it also includes what the parent chain already spent.
Sourcepub fn root_run_id(&self) -> Uuid
pub fn root_run_id(&self) -> Uuid
The top-level run this context belongs to: run_id
itself, or, inside a sub-workflow, the run that started the chain.
The engine stamps it on agent pods as
LABEL_ROOT_RUN_ID; set
the same label on a pod you create yourself (PodRun, JobRun) so
that a retry of the top-level run deletes what a dead attempt left.
§Examples
use ironflow_core::provider::LABEL_ROOT_RUN_ID;
use ironflow_engine::context::WorkflowContext;
fn root_label(ctx: &WorkflowContext) -> (&'static str, String) {
(LABEL_ROOT_RUN_ID, ctx.root_run_id().to_string())
}Sourcepub fn workflow_name(&self) -> &str
pub fn workflow_name(&self) -> &str
The workflow name this run belongs to.
Sourcepub fn total_cost_usd(&self) -> Decimal
pub fn total_cost_usd(&self) -> Decimal
Accumulated cost across all executed steps so far.
Sourcepub fn has_allowed_failure(&self) -> bool
pub fn has_allowed_failure(&self) -> bool
Whether at least one allow_failure step failed during this run.
Sourcepub fn total_duration_ms(&self) -> u64
pub fn total_duration_ms(&self) -> u64
Accumulated duration across all executed steps so far.
Sourcepub fn step_results(&self) -> &[StepResult]
pub fn step_results(&self) -> &[StepResult]
Enriched results of all completed steps in execution order.
Sourcepub async fn payload(&self) -> Result<Value, EngineError>
pub async fn payload(&self) -> Result<Value, EngineError>
Access the payload that triggered this run.
Fetches the run from the store and returns its payload.
§Errors
Returns EngineError::Store if the run is not found.
Sourcepub async fn input<T: DeserializeOwned>(&self) -> Result<T, EngineError>
pub async fn input<T: DeserializeOwned>(&self) -> Result<T, EngineError>
Deserialize the run payload into a typed input struct.
Shorthand for serde_json::from_value(ctx.payload().await?).
§Errors
Returns EngineError::Store if the run is not found, or
EngineError::Serialization if the payload does not match T.
§Examples
use serde::Deserialize;
#[derive(Deserialize)]
struct DeployInput {
environment: String,
dry_run: Option<bool>,
}
let input: DeployInput = ctx.input().await?;Source§impl WorkflowContext
impl WorkflowContext
Sourcepub async fn put_artifact(
&self,
producer: &StepOutput,
name: &str,
content_type: Option<&str>,
content: Vec<u8>,
) -> Result<ArtifactRef, EngineError>
pub async fn put_artifact( &self, producer: &StepOutput, name: &str, content_type: Option<&str>, content: Vec<u8>, ) -> Result<ArtifactRef, EngineError>
Store an in-memory payload as an artifact of the step that produced
producer, and return a handle on it.
The declarative ShellConfig::output
covers shell steps; this covers custom operations and agent steps, which
have no working directory to collect from.
The MIME type is guessed from name unless content_type is set.
§Errors
Returns EngineError::StepConfig when producer does not come from a
recorded step (built by hand, or while planning),
EngineError::ArtifactsUnavailable when no backend is attached,
EngineError::Artifact when the name is invalid or storage fails, and
EngineError::Store when the step already owns that name.
§Examples
use ironflow_engine::context::WorkflowContext;
use ironflow_engine::error::EngineError;
use ironflow_engine::operation::Operation;
let out = ctx.operation("generate", generate).await?;
let summary = ctx
.put_artifact(&out, "summary.json", None, br#"{"ok":true}"#.to_vec())
.await?;
let bytes = ctx.get_artifact(&summary).await?;Sourcepub async fn get_artifact(
&self,
artifact: &ArtifactRef,
) -> Result<Vec<u8>, EngineError>
pub async fn get_artifact( &self, artifact: &ArtifactRef, ) -> Result<Vec<u8>, EngineError>
Read back an artifact produced earlier in this run.
Resolution follows the same rule as a declared input: same run and attempt, steps positioned strictly before the current one, closest producer wins.
§Errors
Returns EngineError::ArtifactNotFound when nothing matches,
EngineError::ArtifactsUnavailable when no backend is attached, and
EngineError::Artifact when the bytes cannot be read.
§Examples
use ironflow_engine::config::ShellConfig;
use ironflow_engine::context::WorkflowContext;
use ironflow_engine::error::EngineError;
let build = ctx.shell("build", ShellConfig::new("./gen").output("report.html")).await?;
let bytes = ctx.get_artifact(&build.artifact("report.html")?).await?;
println!("{} bytes", bytes.len());Source§impl WorkflowContext
impl WorkflowContext
Sourcepub fn on_error(&mut self, name: &str, config: impl Into<StepConfig>)
pub fn on_error(&mut self, name: &str, config: impl Into<StepConfig>)
Register an error handler that fires when any subsequent step fails.
The handler is consumed after firing (fire-once). Multiple handlers can be registered; they fire in registration order.
Error handler execution is best-effort: if a handler fails, the error
is logged but the original step error is preserved. Error handler steps
appear in the run timeline with
Step::is_error_handler
set to true.
§Examples
use ironflow_engine::context::WorkflowContext;
use ironflow_engine::config::ShellConfig;
use ironflow_engine::error::EngineError;
ctx.on_error("cleanup", ShellConfig::new("rm -rf /tmp/build"));
ctx.shell("build", ShellConfig::new("cargo build")).await?;Sourcepub fn clear_error_handlers(&mut self)
pub fn clear_error_handlers(&mut self)
Remove all registered error handlers.
§Examples
use ironflow_engine::context::WorkflowContext;
use ironflow_engine::config::ShellConfig;
use ironflow_engine::error::EngineError;
ctx.on_error("cleanup", ShellConfig::new("rm -rf /tmp/build"));
ctx.shell("build", ShellConfig::new("cargo build")).await?;
ctx.clear_error_handlers();
// cleanup will NOT fire if deploy fails
ctx.shell("deploy", ShellConfig::new("./deploy.sh")).await?;Source§impl WorkflowContext
impl WorkflowContext
Sourcepub async fn agent<C: AgentStep>(
&mut self,
name: &str,
config: C,
) -> Result<C::Answer, EngineError>
pub async fn agent<C: AgentStep>( &mut self, name: &str, config: C, ) -> Result<C::Answer, EngineError>
Execute an agent step.
With output::<T>() on the
config, the step returns the T the agent answered; otherwise it
returns the raw StepOutput. See
AgentStep.
§Errors
Returns EngineError if the agent invocation fails or the store
errors, and EngineError::Serialization if a typed answer does not
match its type.
§Examples
use ironflow_engine::context::WorkflowContext;
use ironflow_engine::config::AgentStepConfig;
use ironflow_engine::error::EngineError;
use schemars::JsonSchema;
use serde::Deserialize;
#[derive(Deserialize, JsonSchema)]
struct Review {
approved: bool,
comments: Vec<String>,
}
let review = ctx
.agent(
"review",
AgentStepConfig::new("Review the code").max_turns(2).output::<Review>(),
)
.await?;
if !review.approved {
println!("{} comments", review.comments.len());
}Source§impl WorkflowContext
impl WorkflowContext
Sourcepub async fn approval(
&mut self,
name: &str,
config: ApprovalConfig,
) -> Result<(), EngineError>
pub async fn approval( &mut self, name: &str, config: ApprovalConfig, ) -> Result<(), EngineError>
Create a human approval gate.
On first execution, records an approval step and returns
EngineError::ApprovalRequired to suspend the run. The engine
transitions the run to AwaitingApproval.
On resume (after a human approved via the API), the approval step
is replayed: it is marked as Completed and execution continues
past it. Multiple approval gates in the same handler work – each
one pauses and resumes independently.
When the config carries an SLA
(ApprovalConfig::with_deadline,
or the legacy with_timeout_seconds), the deadline is persisted on the
step so the API server’s escalator can apply the configured
EscalationPolicy when it fires. The
timer is cleared as soon as the gate resolves.
A StepInterceptor wired into the
context resolves the gate inline instead of suspending: the step is
recorded, then completed or rejected without waiting for a human. This
is what crate::testing::TestEngine uses to run gated handlers end to
end.
When the config carries Approvers
(ApprovalConfig::requiring),
they are recorded on the step as an
ApprovalRequirement when the gate
opens. That record stays the source of truth on replay and resume: the
approvers the handler computes on a later execution are ignored. A config
without approvers stores no requirement.
§Errors
Returns EngineError::ApprovalRequired to pause the run on
first execution. Returns EngineError::ApprovalRejected when an
interceptor refuses the gate. Returns other EngineError variants on
store failures.
§Examples
use ironflow_engine::context::WorkflowContext;
use ironflow_engine::config::{ApprovalConfig, Approvers};
use ironflow_engine::error::EngineError;
use serde::Deserialize;
#[derive(Deserialize)]
struct Payment {
amount: u64,
}
let payment: Payment = ctx.input().await?;
let approvers = match payment.amount {
a if a > 10_000 => Approvers::at_least(2).from_groups(["finance"]).because("amount > 10k"),
_ => Approvers::any(),
};
ctx.approval("payment-gate", ApprovalConfig::new("Release the payment?").requiring(approvers))
.await?;
// Execution continues here after approvalSource§impl WorkflowContext
impl WorkflowContext
Sourcepub async fn when<T, F>(
&mut self,
label: &str,
predicate: F,
) -> Result<bool, EngineError>
pub async fn when<T, F>( &mut self, label: &str, predicate: F, ) -> Result<bool, EngineError>
Evaluate a named branch condition against the typed run input.
The run payload is deserialized into T, exactly like
input, and predicate decides the branch on it. label
is a human-readable name for the branch, shown in the plan; it is never
parsed nor evaluated. In plan mode the result is also recorded as
ConditionResult::Evaluated on the next planned step, so the operator
sees which branch the plan followed and why.
§Errors
Returns EngineError::Store when the run payload cannot be read, and
EngineError::Serialization when it does not match T: a misspelled
field or variant fails the branch instead of silently taking the other
one.
§Examples
use ironflow_engine::context::WorkflowContext;
use ironflow_engine::config::ShellConfig;
use ironflow_engine::error::EngineError;
use serde::Deserialize;
#[derive(Deserialize, PartialEq)]
#[serde(rename_all = "lowercase")]
enum Env {
Prod,
Staging,
}
#[derive(Deserialize)]
struct DeployInput {
env: Env,
}
if ctx.when("production run", |i: &DeployInput| i.env == Env::Prod).await? {
ctx.shell("deploy-prod", ShellConfig::new("./deploy prod")).await?;
} else {
ctx.skip("deploy-prod", "not a production run").await?;
}Sourcepub fn when_dynamic(&mut self, label: &str, value: bool) -> bool
pub fn when_dynamic(&mut self, label: &str, value: bool) -> bool
Record a branch condition whose value depends on a previous step’s output.
Returns value unchanged; label names the branch in the plan. Under
planning, step outputs are synthetic,
so the condition is recorded as ConditionResult::Unevaluable: the
plan still follows the branch the synthetic output produces, and the
operator is told the other branch may run instead.
§Examples
use ironflow_engine::context::WorkflowContext;
use ironflow_engine::config::ShellConfig;
use ironflow_engine::error::EngineError;
let build = ctx.shell("build", ShellConfig::new("cargo build")).await?;
if ctx.when_dynamic("build succeeded", build.is_success()) {
ctx.shell("deploy", ShellConfig::new("./deploy")).await?;
}Source§impl WorkflowContext
impl WorkflowContext
Sourcepub async fn decision<T: DecisionAnswers>(
&mut self,
name: &str,
config: DecisionConfig<T>,
) -> Result<T, EngineError>
pub async fn decision<T: DecisionAnswers>( &mut self, name: &str, config: DecisionConfig<T>, ) -> Result<T, EngineError>
Execute a typed machine-decision step (System One / Jev).
The questions come from T, set with
DecisionConfig::answers, and the answers are returned as a T. See
crate::decision. When escalate_below is set and any answer falls
below it, the run suspends with EngineError::ApprovalRequired and
replays the stored answers on resume without re-calling the provider.
While planning no provider is called and no answer exists, so reading
them fails with EngineError::Decision and the plan stops at this
step.
§Errors
EngineError::NoDecisionProvider, EngineError::ApprovalRequired,
EngineError::Operation, or EngineError::Decision when an answer
does not fit T.
§Examples
use ironflow_engine::config::DecisionConfig;
use ironflow_engine::context::WorkflowContext;
use ironflow_engine::decision::DecisionAnswers;
use ironflow_engine::error::EngineError;
#[derive(DecisionAnswers)]
struct Urgency {
#[noul("Does this convey urgency?")]
urgent: f64,
}
let urgency = ctx
.decision("urgency", DecisionConfig::new("Payouts fail since 3 days").answers::<Urgency>())
.await?;
if urgency.urgent > 0.8 {
// page someone
}A config without questions is not a decision:
ctx.decision("urgency", DecisionConfig::new("state")).await?;Source§impl WorkflowContext
impl WorkflowContext
Sourcepub async fn delay(
&mut self,
name: &str,
config: DelayConfig,
) -> Result<(), EngineError>
pub async fn delay( &mut self, name: &str, config: DelayConfig, ) -> Result<(), EngineError>
Execute a delay (timed pause) step.
A zero-duration delay completes immediately. Otherwise, the
delay step is marked completed and the method returns
EngineError::DelaySleeping so the engine transitions the
run to Sleeping.
On resume (after the worker picks up the re-queued run), the delay step is replayed as completed via the replay mechanism.
§Errors
Returns EngineError::DelaySleeping to suspend the run.
§Examples
use ironflow_engine::context::WorkflowContext;
use ironflow_engine::config::delay::DelayConfig;
use ironflow_engine::error::EngineError;
ctx.delay("cooldown", DelayConfig::from_secs(300)).await?;Source§impl WorkflowContext
impl WorkflowContext
Sourcepub async fn http(
&mut self,
name: &str,
config: HttpConfig,
) -> Result<StepOutput, EngineError>
pub async fn http( &mut self, name: &str, config: HttpConfig, ) -> Result<StepOutput, EngineError>
Execute an HTTP step.
§Errors
Returns EngineError if the request fails or the store errors.
§Examples
use ironflow_engine::context::WorkflowContext;
use ironflow_engine::config::HttpConfig;
use ironflow_engine::error::EngineError;
let resp = ctx.http("health", HttpConfig::get("https://api.example.com/health")).await?;
println!("status: {:?}, body: {}", resp.status(), resp.body());Source§impl WorkflowContext
impl WorkflowContext
Sourcepub async fn human_input<T: DeserializeOwned + JsonSchema>(
&mut self,
name: &str,
config: HumanInputConfig,
) -> Result<T, EngineError>
pub async fn human_input<T: DeserializeOwned + JsonSchema>( &mut self, name: &str, config: HumanInputConfig, ) -> Result<T, EngineError>
Ask a human for a typed answer and suspend the run until it is given.
On first execution, records a human input step carrying the JSON schema
of T and returns EngineError::HumanInputRequired to suspend the
run. The engine transitions the run to AwaitingApproval. The answer is
posted to POST /api/v1/runs/{id}/steps/{step_id}/input, validated
against the schema, and the run resumes.
On resume, the step is replayed and the stored answer is deserialized
into T. A rejected input (POST .../steps/{step_id}/reject) returns
EngineError::HumanInputRejected so the handler decides what happens
next. An answer given in an earlier attempt is carried over to a retry.
The config reuses the approval gate machinery: deadline, escalation
policy, assignee and the Approvers allowed to answer.
While planning, the step is recorded and never suspends: T is
deserialized from {} when it accepts that (for example with
#[serde(default)]).
§Errors
Returns EngineError::HumanInputRequired to pause the run until an
answer is given. Returns EngineError::HumanInputRejected when the
input was rejected. Returns EngineError::StepConfig when the stored
answer does not match T, or while planning when T cannot be built
from {}. Returns other EngineError variants on store failures.
§Examples
use ironflow_engine::config::HumanInputConfig;
use ironflow_engine::context::WorkflowContext;
use ironflow_engine::error::EngineError;
use schemars::JsonSchema;
use serde::Deserialize;
#[derive(Deserialize, JsonSchema)]
struct Answers {
answers: Vec<String>,
}
let answers: Answers = ctx
.human_input("clarify", HumanInputConfig::new("Answer the clarification questions"))
.await?;
assert!(answers.answers.len() < 100);Source§impl WorkflowContext
impl WorkflowContext
Sourcepub async fn operation(
&mut self,
name: &str,
op: &dyn Operation,
) -> Result<StepOutput, EngineError>
pub async fn operation( &mut self, name: &str, op: &dyn Operation, ) -> Result<StepOutput, EngineError>
Execute a custom operation step.
Runs a user-defined Operation with full step lifecycle management:
creates the step record, transitions to Running, executes the operation,
persists the output and duration, and marks the step Completed or Failed.
The operation’s kind() is stored as
StepKind::Custom.
On resume, a step that already completed in a prior execution of the
current attempt is replayed from the store instead of calling
Operation::execute again.
§Errors
Returns EngineError if the operation fails or the store errors.
§Examples
use async_trait::async_trait;
use ironflow_engine::context::WorkflowContext;
use ironflow_engine::operation::{Operation, OperationContext};
use ironflow_core::error::OperationError;
use ironflow_engine::error::EngineError;
use serde_json::{Value, json};
struct MyOp;
#[async_trait]
impl Operation for MyOp {
fn kind(&self) -> &str { "my-service" }
async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
Ok(json!({"ok": true}))
}
}
let result = ctx.operation("call-service", &MyOp).await?;
println!("output: {}", result.output);Source§impl WorkflowContext
impl WorkflowContext
Sourcepub async fn parallel(
&mut self,
steps: Vec<(&str, StepConfig)>,
fail_fast: bool,
) -> Result<Vec<ParallelStepResult>, EngineError>
pub async fn parallel( &mut self, steps: Vec<(&str, StepConfig)>, fail_fast: bool, ) -> Result<Vec<ParallelStepResult>, EngineError>
Execute multiple steps concurrently (wait-all model).
All steps in the batch execute in parallel via tokio::JoinSet.
Each step is recorded with the same position (execution wave).
Dependencies on previous steps are recorded automatically.
When fail_fast is true, remaining steps are aborted on the first
failure. When false, all steps run to completion and the first
error is returned.
Every step of a wave must have its own name: the name identifies the
step in the run timeline, in its artifact handles, on resume and in the
ironflow.io/step pod label. Two steps sharing that label would let
the K8s ephemeral provider delete one step’s pod when starting the other.
On resume, each step of the wave that already completed in a prior execution of the current attempt is replayed from the store; only the other steps of the wave are launched again.
§Errors
Returns EngineError::StepConfig if two steps of the wave share a
name, before anything runs. Returns EngineError if any step fails.
§Examples
use ironflow_engine::context::WorkflowContext;
use ironflow_engine::config::{StepConfig, ShellConfig};
use ironflow_engine::error::EngineError;
let results = ctx.parallel(
vec![
("test-unit", StepConfig::Shell(ShellConfig::new("cargo test --lib"))),
("lint", StepConfig::Shell(ShellConfig::new("cargo clippy"))),
],
true,
).await?;
for r in &results {
println!("{}: {:?}", r.name, r.output.output);
}Source§impl WorkflowContext
impl WorkflowContext
Sourcepub async fn shell(
&mut self,
name: &str,
config: ShellConfig,
) -> Result<StepOutput, EngineError>
pub async fn shell( &mut self, name: &str, config: ShellConfig, ) -> Result<StepOutput, EngineError>
Execute a shell step.
Creates the step record, runs the command, persists the result, and returns the output for use in subsequent steps.
§Errors
Returns EngineError if the command fails or the store errors.
§Examples
use ironflow_engine::context::WorkflowContext;
use ironflow_engine::config::ShellConfig;
use ironflow_engine::error::EngineError;
let files = ctx.shell("list", ShellConfig::new("ls -la")).await?;
println!("stdout: {}", files.stdout());Source§impl WorkflowContext
impl WorkflowContext
Sourcepub async fn skip(
&mut self,
name: &str,
reason: &str,
) -> Result<(), EngineError>
pub async fn skip( &mut self, name: &str, reason: &str, ) -> Result<(), EngineError>
Record a step as explicitly skipped.
Use this inside an if/else branch when a step should not execute
but must still appear in the DAG and timeline with its reason.
The step is created directly in StepStatus::Skipped state and the
reason is stored in the output as {"reason": "..."}.
On resume, a skip already recorded in a prior execution of the current
attempt is replayed instead of creating a second Skipped step.
§Errors
Returns EngineError if the store fails.
§Examples
use ironflow_engine::context::WorkflowContext;
use ironflow_engine::error::EngineError;
let tests_passed = false;
if tests_passed {
// ctx.shell("deploy", ...).await?;
} else {
ctx.skip("deploy", "tests failed").await?;
}Source§impl WorkflowContext
impl WorkflowContext
Sourcepub async fn workflow<W: TypedWorkflow>(
&mut self,
handler: &W,
input: W::Input,
) -> Result<SubWorkflowOutput, EngineError>
pub async fn workflow<W: TypedWorkflow>( &mut self, handler: &W, input: W::Input, ) -> Result<SubWorkflowOutput, EngineError>
Execute a sub-workflow step.
Creates a child run of handler whose payload is input, executes it
with its own steps and lifecycle, and returns its run ID and aggregated
metrics. The child declares its input type through TypedWorkflow,
so only a W::Input is accepted.
Requires the context to be created with
with_handler_resolver.
§Errors
Returns EngineError::InvalidWorkflow if no handler is registered
with the given name, or if no handler resolver is available, and
EngineError::Serialization if input cannot be serialized.
§Examples
use ironflow_engine::context::WorkflowContext;
use ironflow_engine::error::EngineError;
use ironflow_engine::handler::{HandlerFuture, TypedWorkflow, WorkflowHandler};
use serde::{Deserialize, Serialize};
#[derive(Serialize, Deserialize)]
struct CollectInput {
scope: String,
}
struct Collect;
impl WorkflowHandler for Collect {
fn name(&self) -> &str { "collect" }
fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
Box::pin(async move { Ok(()) })
}
}
impl TypedWorkflow for Collect {
type Input = CollectInput;
}
let child = ctx.workflow(&Collect, CollectInput { scope: "system".to_string() }).await?;
let steps = ctx.store().list_steps(child.run_id()).await?;Any other input type is a compile error:
ctx.workflow(&Collect, serde_json::json!({"scope": "system"})).await?;Sourcepub async fn workflow_dyn(
&mut self,
handler: &dyn WorkflowHandler,
payload: Value,
) -> Result<SubWorkflowOutput, EngineError>
👎Deprecated: implement TypedWorkflow on the child and call workflow: its payload is then checked at compile time
pub async fn workflow_dyn( &mut self, handler: &dyn WorkflowHandler, payload: Value, ) -> Result<SubWorkflowOutput, EngineError>
implement TypedWorkflow on the child and call workflow: its payload is then checked at compile time
Execute a sub-workflow step whose child is only known at run time.
Same as workflow, without the compile-time check of
the payload: the child must deserialize payload itself.
§Errors
Same as workflow.
§Examples
use ironflow_engine::context::WorkflowContext;
use ironflow_engine::error::EngineError;
use ironflow_engine::handler::WorkflowHandler;
use serde_json::json;
let result = ctx.workflow_dyn(child, json!({"scope": "system"})).await?;
println!("child run {}", result.run_id());