Ironflow
Ironflow is a workflow orchestration platform where workflows are imperative Rust code executed by background workers, with persistence, cost tracking, and human approval gates.
Why Ironflow?
- Workflows are Rust code. No YAML, no DSL. Full type safety, IDE support, and compile-time checks.
- Persistent execution. Every step is tracked in a database. Runs survive process restarts.
- Human-in-the-loop. Approval gates pause a run until someone approves or rejects it.
- Cost tracking. Every agent call is metered. Set per-run and monthly budgets.
- Scalable. Add more workers to increase throughput. Workers poll the API for pending runs.
Quick links
- Getting Started – install, configure, run
- Concepts – understand the building blocks
- Guides – step-by-step tutorials
- Architecture – how the pieces fit together
- API Reference (docs.rs) – generated from rustdoc
Installation
Prerequisites
- Rust 1.94+ (see
rust-versionin Cargo.toml) - A running PostgreSQL instance (for production; in-memory store available for development)
Add dependencies
Add the crates you need to your Cargo.toml:
[dependencies]
ironflow-engine = "0.1" # Workflow handler, context, engine
ironflow-api = "0.1" # REST API server
ironflow-worker = "0.1" # Background worker
ironflow-store = "0.1" # Storage backends
ironflow-core = "0.1" # Shell, agent providers
Minimal project structure
A typical Ironflow project has three parts:
- A library crate with your workflow handlers
- A server binary that exposes the API and serves the dashboard
- A worker binary that executes workflows
my-project/
├── src/
│ └── lib.rs # Your workflow handlers
├── src/bin/
│ ├── server.rs # API server
│ └── worker.rs # Background worker
└── Cargo.toml
See the example server and example worker for complete working code.
Running the Server
The API server exposes a REST API for managing workflows, runs, and steps. It also serves the web dashboard.
Example server
The repository includes a complete example server:
//! ironflow API server example.
//!
//! ```sh
//! cargo run -p ironflow-example-server
//! ```
//!
//! The dashboard is served automatically via the `dashboard` feature in `ironflow-api`.
//!
//! Environment:
//! - `IRONFLOW_ENV` (`production` or `development`, default: development)
//! - `DATABASE_URL` (required in production)
//! - `JWT_SECRET` (required in production, default: dev secret)
//! - `WORKER_TOKEN` (required in production, default: dev token)
//! - `PORT` (default: 3000)
//! - `DASHBOARD_DIR` (optional: overrides the embedded dashboard with a filesystem path)
//! - `ALLOWED_ORIGINS` (comma-separated list; omit to allow same-origin only)
//! - `WEBHOOK_URL` (optional: outbound webhook for run events)
//! - `ARTIFACTS_DIR` (optional: filesystem root for step artifacts; unset
//! leaves artifacts disabled and the artifact routes answer `501`)
//! - `ARTIFACT_MAX_BYTES` (optional: per-artifact size limit, default 100 MiB)
//! - `IRONFLOW_DEFAULT_RUN_MAX_COST_USD` (optional: default per-run cost cap in
//! USD, applied when neither the run creation request nor the workflow
//! handler declares one; unset means no cap)
//! - `IRONFLOW_MONTHLY_COST_LIMIT_USD` (optional: global cost quota for the
//! current calendar month in UTC; beyond it, creating a run returns
//! `429 MONTHLY_BUDGET_EXCEEDED` while in-flight runs continue)
//! - `PURGE_MAX_AGE_DAYS`, `PURGE_MAX_RUNS_PER_WORKFLOW`, `PURGE_DRY_RUN`,
//! `PURGE_INTERVAL_SECS`, `PROVIDER_ACCOUNT_USAGE_RETENTION_DAYS` and
//! `SIGNAL_RETENTION_DAYS` (optional: retention of the purger, see
//! `ServerConfig`)
//! - `IRONFLOW_SEED` (optional: when set to any value, seeds development data
//! at startup -- users, runs, steps, API keys)
use std::process;
use std::sync::Arc;
use axum::http::header::{AUTHORIZATION, CONTENT_TYPE};
use axum::http::{HeaderValue, Method};
use tokio::net::TcpListener;
use tower_http::cors::CorsLayer;
use tracing::{info, warn};
use tracing_subscriber::EnvFilter;
use ironflow_api::config::ServerConfig;
use ironflow_api::purger::RunPurger;
use ironflow_api::routes::{RouterConfig, create_router};
use ironflow_api::sse::SseBroadcaster;
use ironflow_api::state::AppState;
use ironflow_artifacts::blob_store::BlobStore;
use ironflow_artifacts::local::LocalBlobStore;
use ironflow_auth::jwt::JwtConfig;
use ironflow_core::providers::claude::ClaudeCodeProvider;
use ironflow_engine::artifact::DirectArtifactSink;
use ironflow_engine::budget::BudgetConfig;
use ironflow_engine::engine::Engine;
use ironflow_engine::notify::{Event, WebhookSubscriber, WorkflowEventBus};
use ironflow_store::crypto::{KeyRing, SECRET_KEYS_ENV};
use ironflow_store::memory::InMemoryStore;
use ironflow_store::store::Store;
use xtask::seed::{SeedOptions, seed_store};
#[tokio::main]
async fn main() {
tracing_subscriber::fmt()
.with_env_filter(
EnvFilter::try_from_default_env()
.unwrap_or_else(|_| "info,ironflow=debug".parse().expect("valid filter")),
)
.init();
let config = ServerConfig::from_env().unwrap_or_else(|e| {
eprintln!("{e}");
process::exit(1);
});
let mut store = InMemoryStore::new();
let key_ring = KeyRing::from_env().unwrap_or_else(|e| {
eprintln!("invalid secret key configuration: {e}");
process::exit(1);
});
let has_key_ring = key_ring.is_some();
match key_ring {
Some(ring) => {
info!(
active_version = ring.active_version(),
configured_versions = ?ring.versions(),
"secret store enabled"
);
store.set_key_ring(ring);
}
None => {
info!("{SECRET_KEYS_ENV} not set, secret store disabled");
}
}
let store: Arc<dyn Store> = Arc::new(store);
// A secret encrypted with a key that is no longer configured is
// unreadable. Fail here rather than at the first workflow that needs it.
if has_key_ring {
let status = store.secret_key_status().await.unwrap_or_else(|e| {
eprintln!("cannot read secret key versions: {e}");
process::exit(1);
});
if !status.is_consistent() {
let missing: Vec<String> = status.missing.iter().map(|v| v.to_string()).collect();
eprintln!(
"secret key versions present in database but missing from configuration: {}\n\
set {SECRET_KEYS_ENV} to include them, or rotate before removing a key",
missing.join(", ")
);
process::exit(1);
}
}
if std::env::var("IRONFLOW_SEED").is_ok() {
info!("IRONFLOW_SEED set, seeding development data...");
let seed_opts = SeedOptions {
force: false,
artifacts_dir: config.artifacts_dir.clone(),
};
seed_store(&*store, &seed_opts).await.unwrap_or_else(|e| {
warn!("seed skipped: {e}");
});
}
let provider = Arc::new(ClaudeCodeProvider::new());
let jwt_config = Arc::new(JwtConfig {
secret: config.jwt_secret.clone(),
access_token_ttl_secs: 900,
refresh_token_ttl_secs: 604800,
cookie_domain: None,
cookie_secure: config.is_production,
});
let budget = BudgetConfig::from_env();
info!(
default_run_max_cost_usd = ?budget.default_run_max_cost_usd,
monthly_cost_limit_usd = ?budget.monthly_cost_limit_usd,
"cost guardrails loaded"
);
let mut engine = Engine::new(store.clone(), provider).with_budget_config(budget);
ironflow_workflows::register_all(&mut engine).expect("failed to register workflows");
// Artifacts stay off until a storage root is configured. The API and any
// in-process run then share the same backend, so a file a step produces is
// downloadable from the same server that stored it.
let blob_store: Option<Arc<dyn BlobStore>> = config.artifacts_dir.as_ref().map(|dir| {
info!(
dir = %dir.display(),
max_bytes = config.artifact_max_bytes,
"artifact storage enabled"
);
Arc::new(LocalBlobStore::new(dir).max_bytes(config.artifact_max_bytes))
as Arc<dyn BlobStore>
});
if let Some(ref blob) = blob_store {
engine.set_artifact_sink(Arc::new(DirectArtifactSink::new(
blob.clone(),
store.clone(),
)));
}
if let Some(ref webhook_url) = config.webhook_url {
info!(url = %webhook_url, "registering webhook subscriber");
engine.subscribe(
WebhookSubscriber::new(webhook_url),
&[Event::RUN_STATUS_CHANGED, Event::STEP_FAILED],
);
}
let sse_broadcaster = SseBroadcaster::new();
let event_sender = sse_broadcaster.sender();
engine.subscribe(sse_broadcaster, Event::ALL);
let event_bus = WorkflowEventBus::new();
engine.set_event_bus(event_bus.clone());
let engine = Arc::new(engine);
let cors = build_cors(&config);
let mut state = AppState::new(
store.clone(),
engine.clone(),
jwt_config,
config.worker_token.clone(),
event_sender,
)
.with_event_bus(event_bus);
if let Some(blob) = blob_store {
state = state.with_blob_store(blob);
}
let shutdown = state.spawn_background_tasks().await;
tokio::spawn(
RunPurger::from_config(store.clone(), &config)
.with_blob_store(state.blob_store.clone())
.run(shutdown.clone()),
);
let router_config = RouterConfig {
dashboard_dir: config.dashboard_dir.clone(),
rate_limit_auth: config.rate_limit_auth,
rate_limit_general: config.rate_limit_general,
};
let app = create_router(state, router_config)
.layer(cors)
.into_make_service();
let addr = format!("0.0.0.0:{}", config.port);
let listener = TcpListener::bind(&addr).await.expect("bind address");
info!("==============================================");
info!(" ironflow server on http://{addr}");
info!(
" environment: {}",
if config.is_production {
"production"
} else {
"development"
}
);
info!("==============================================");
axum::serve(listener, app)
.with_graceful_shutdown(async move {
tokio::signal::ctrl_c().await.expect("ctrl+c handler");
info!("shutting down...");
shutdown.cancel();
})
.await
.expect("serve");
}
/// Build CORS layer from config.
///
/// - If `allowed_origins` is set: only those origins are permitted (comma-separated).
/// - If unset: no extra origins are allowed (same-origin only).
///
/// Credentials (cookies) are always allowed so JWT cookies work cross-origin.
fn build_cors(config: &ServerConfig) -> CorsLayer {
let methods = vec![Method::GET, Method::POST, Method::PUT, Method::DELETE];
let headers = vec![AUTHORIZATION, CONTENT_TYPE];
match config.allowed_origins {
Some(ref raw) => {
let origins: Vec<HeaderValue> = raw
.split(',')
.map(str::trim)
.filter(|s| !s.is_empty())
.filter_map(|s| match s.parse::<HeaderValue>() {
Ok(v) => Some(v),
Err(err) => {
warn!(origin = s, %err, "ignoring invalid CORS origin");
None
}
})
.collect();
info!(?origins, "CORS: allowing configured origins");
CorsLayer::new()
.allow_origin(origins)
.allow_methods(methods)
.allow_headers(headers)
.allow_credentials(true)
}
None => {
info!("CORS: no ALLOWED_ORIGINS set, same-origin only");
CorsLayer::new()
.allow_methods(methods)
.allow_headers(headers)
}
}
}
Environment variables
| Variable | Default | Description |
|---|---|---|
IRONFLOW_ENV | development | production or development |
DATABASE_URL | – | PostgreSQL URL (required in production) |
JWT_SECRET | dev secret | JWT signing key (in production: mandatory, >= 32 bytes, must not start with ironflow-dev-) |
WORKER_TOKEN | dev token | Shared secret for worker auth (in production: mandatory, >= 32 bytes, must not start with ironflow-dev-) |
PORT | 3000 | HTTP listen port |
ALLOWED_ORIGINS | same-origin | Comma-separated CORS origins |
ARTIFACTS_DIR | – | Filesystem root for step artifacts |
PURGE_MAX_AGE_DAYS | 90 | Terminal runs older than this are purged |
PURGE_MAX_RUNS_PER_WORKFLOW | 1000 | Terminal runs kept per workflow |
PURGE_DRY_RUN | false | Log what would be purged without deleting (also disables usage and signal purging) |
PURGE_INTERVAL_SECS | 86400 | Seconds between purger ticks (min 60) |
PROVIDER_ACCOUNT_USAGE_RETENTION_DAYS | 30 | Days of Provider Account usage history kept (min 1) |
SIGNAL_RETENTION_DAYS | 7 | Days received signals are kept before the purger deletes them (min 1) |
Retention
The example server starts a RunPurger built from these variables: it purges old runs,
Provider Account usage history and received signals. A custom server must start it
itself with RunPurger::from_config, otherwise these variables are ignored:
let shutdown = state.spawn_background_tasks().await;
tokio::spawn(
RunPurger::from_config(store.clone(), &config)
.with_blob_store(state.blob_store.clone())
.run(shutdown.clone()),
);
Running
cargo run -p ironflow-example-server
The server starts on http://localhost:3000. The dashboard is available at the root URL.
Running a Worker
Workers poll the API for pending runs, acquire leases, and execute workflow handlers.
Example worker
//! ironflow worker example.
//!
//! ```sh
//! cargo run -p ironflow-example-worker
//! ```
//!
//! Environment:
//! - `API_URL` (default: http://localhost:3000)
//! - `WORKER_TOKEN` (default: dev token)
//! - `CONCURRENCY` (default: 2)
//! - `POLL_INTERVAL_SECS` (default: 2)
use std::env;
use std::sync::Arc;
use std::time::Duration;
use tracing::info;
use tracing_subscriber::EnvFilter;
use ironflow_core::providers::claude::ClaudeCodeProvider;
use ironflow_worker::WorkerBuilder;
use ironflow_workflows::handlers;
#[tokio::main]
async fn main() {
tracing_subscriber::fmt()
.with_env_filter(
EnvFilter::try_from_default_env()
.unwrap_or_else(|_| "info,ironflow=debug".parse().expect("valid filter")),
)
.init();
let api_url = env::var("API_URL").unwrap_or_else(|_| "http://localhost:3000".to_string());
let worker_token =
env::var("WORKER_TOKEN").unwrap_or_else(|_| "ironflow-dev-worker-token".to_string());
let concurrency: usize = env::var("CONCURRENCY")
.ok()
.and_then(|c| c.parse().ok())
.unwrap_or(2);
let poll_interval: u64 = env::var("POLL_INTERVAL_SECS")
.ok()
.and_then(|p| p.parse().ok())
.unwrap_or(2);
let mut builder = WorkerBuilder::new(&api_url, &worker_token)
.provider(Arc::new(ClaudeCodeProvider::new()))
.concurrency(concurrency)
.poll_interval(Duration::from_secs(poll_interval));
// Same list as the server: one source of truth for both binaries.
for handler in handlers() {
builder = builder.register(handler);
}
let worker = builder.build().expect("failed to build worker");
info!("==============================================");
info!(" ironflow worker");
info!(" API: {api_url}");
info!(" Concurrency: {concurrency}");
info!("==============================================");
if let Err(e) = worker.run().await {
tracing::error!("worker error: {e}");
}
}
Environment variables
| Variable | Default | Description |
|---|---|---|
API_URL | http://localhost:3000 | Address of the API server |
WORKER_TOKEN | dev token | Shared secret matching the server |
CONCURRENCY | 2 | Number of parallel runs |
POLL_INTERVAL_SECS | 2 | Seconds between polls |
Running
cargo run -p ironflow-example-worker
Scaling
To increase throughput, start multiple workers. Each worker polls independently and acquires leases on runs, so no coordination is needed beyond the API server.
WorkflowHandler
A WorkflowHandler is the core abstraction in Ironflow. It defines a named workflow as imperative Rust code.
The trait
pub trait WorkflowHandler: Send + Sync {
fn name(&self) -> &str;
fn description(&self) -> &str;
fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a>;
// Optional methods
fn category(&self) -> Option<&str> { None }
fn input_schema(&self) -> Option<Value> { None }
fn default_labels(&self) -> HashMap<String, String> { HashMap::new() }
fn source_code(&self) -> Option<&str> { None }
}
Example: a greeting workflow
use std::collections::HashMap;
use ironflow_engine::config::ShellConfig;
use ironflow_engine::context::WorkflowContext;
use ironflow_engine::handler::{HandlerFuture, WorkflowHandler, input_schema_for};
use schemars::JsonSchema;
use serde::Deserialize;
use serde_json::Value;
/// Input payload for the greeting workflow.
///
/// Derives [`JsonSchema`] so the dashboard can render a dynamic form.
#[derive(Deserialize, JsonSchema)]
struct GreetingInput {
/// Person to greet.
name: String,
/// Greeting language (en, fr, es).
#[serde(default = "default_language")]
language: String,
/// Number of times to repeat the greeting.
#[serde(default = "default_repeat")]
repeat: u32,
/// Whether to output in uppercase.
#[serde(default)]
uppercase: bool,
}
fn default_language() -> String {
"en".to_string()
}
fn default_repeat() -> u32 {
1
}
pub struct Greeting;
impl WorkflowHandler for Greeting {
fn name(&self) -> &str {
"greeting"
}
fn category(&self) -> Option<&str> {
Some("examples")
}
fn input_schema(&self) -> Option<Value> {
Some(input_schema_for::<GreetingInput>())
}
fn default_labels(&self) -> HashMap<String, String> {
HashMap::from([("project".to_string(), "ironflow".to_string())])
}
fn description(&self) -> &str {
"A demo workflow that greets someone. \
Shows how input_schema generates a dynamic form in the dashboard."
}
fn source_code(&self) -> Option<&str> {
Some(include_str!("greeting.rs"))
}
fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
Box::pin(async move {
let input: GreetingInput = ctx.input().await?;
let greeting = match input.language.as_str() {
"fr" => format!("Bonjour, {} !", input.name),
"es" => format!("Hola, {}!", input.name),
_ => format!("Hello, {}!", input.name),
};
let mut message = (0..input.repeat)
.map(|_| greeting.as_str())
.collect::<Vec<_>>()
.join("\n");
if input.uppercase {
message = message.to_uppercase();
}
ctx.shell("greet", ShellConfig::new(&format!("echo '{message}'")))
.await?;
Ok(())
})
}
}
Key points
name()must be unique across all registered handlers. It identifies the workflow in the API and the database.execute()receives aWorkflowContextto create steps. Steps are persisted as they complete.input_schema()returns a JSON Schema derived from a#[derive(JsonSchema)]struct. The dashboard renders it as a dynamic form.source_code()optionally embeds the handler source for display in the dashboard.sub_workflows()lists the handlers this one calls throughctx.workflow, for the call graph. Build it from the handlers,sub_workflow_names(&[&Collect]), never from hand-written names.
Typed input for sub-workflows
A handler called as a sub-workflow declares its input type with TypedWorkflow.
ctx.workflow(&Collect, CollectInput { .. }) then accepts nothing else, and
input_schema() is derived from the same type:
#[derive(Serialize, Deserialize, JsonSchema)]
struct CollectInput {
host: String,
}
impl WorkflowHandler for Collect {
fn name(&self) -> &str { "collect" }
fn input_schema(&self) -> Option<Value> { Self::typed_input_schema() }
fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> { /* .. */ }
}
impl TypedWorkflow for Collect {
type Input = CollectInput;
}
// In the parent:
let child = ctx.workflow(&Collect, CollectInput { host: "db-1".into() }).await?;
let steps = ctx.store().list_steps(child.run_id()).await?;
A child without input uses type Input = ();.
Registration
Handlers are registered in the Engine before starting the server or worker:
let mut engine = Engine::new(store, provider);
engine.register(Box::new(Greeting))?;
See Writing a Workflow for a step-by-step guide.
Steps
A Step is an atomic unit of work within a Run. Each step is persisted in the database with its input, output, status, cost, duration, and token counts.
Step kinds
| Kind | Method | Description |
|---|---|---|
| Shell | ctx.shell() | Execute a shell command |
| Http | ctx.http() | Make an HTTP request |
| Agent | ctx.agent() | Call an AI agent (Claude, OpenAI, etc.) |
| Approval | ctx.approval() | Pause for human approval |
| HumanInput | ctx.human_input() | Pause until a human submits a typed answer (Human Input) |
| Signal | ctx.wait_for_signal() | Pause until an external signal arrives or a timeout elapses (Signals) |
| Decision | ctx.decision() | Make a typed machine decision (System One / Jev) |
| Workflow | ctx.workflow() | Start a sub-workflow |
| Custom | ctx.operation() | Run a custom Operation |
Shell steps
let output = ctx.shell("build", ShellConfig::new("cargo build")).await?;
if output.is_success() {
// continue
}
HTTP steps
let response = ctx.http("fetch-data", HttpConfig::get("https://api.example.com/data")).await?;
A shell step reads back through stdout(), stderr() and exit_code(), an
HTTP step through status() and body(). Files a shell step declares with
.output("target/*.log") are handed to later steps through a handle:
build.artifact("build.log")?, passed to ShellConfig::input(&handle) or
ctx.get_artifact(&handle). A name the step did not declare fails with
EngineError::ArtifactNotDeclared.
Agent steps
let result = ctx.agent("analyze", AgentStepConfig::new("Analyze this log file")).await?;
println!("{}", result.text());
#[derive(Deserialize, JsonSchema)]
struct Verdict { approved: bool }
// With `.output::<T>()` the step returns the `T` itself.
let verdict = ctx
.agent("review", AgentStepConfig::new("Review the diff").max_turns(2).output::<Verdict>())
.await?;
Tools are an enum: .allow_tool(Tool::Bash), Tool::Custom("mcp__server__tool".into())
for anything else. Tools and structured output are mutually exclusive.
On an HTTP provider, tools are registered on the worker in named profiles
(HttpAgentProvider::with_tool_profile). Each profile is a ToolProfile constant,
const BUG: ToolProfile = ToolProfile::new("bug");, shared by the worker and the
handlers, so a misspelled profile does not compile. A step picks one with
.tool_profile(BUG) and sees only its tools; without a profile it gets the with_tools registry, or none.
An unknown profile fails the step, a Claude CLI provider refuses any profile, and the
run logs name the profile with the tools it exposed.
Sub-workflow steps
A child declares its input type with TypedWorkflow and the parent passes that
type; see WorkflowHandler.
let child = ctx.workflow(&Collect, CollectInput { host: "db-1".into() }).await?;
println!("child run {}", child.run_id());
Approval steps
ctx.approval("prod-gate", ApprovalConfig::new("Deploy to production?")).await?;
The run suspends on AwaitingApproval until a human answers. A gate can also
carry an SLA: a deadline persisted on the step and an escalation policy applied
when it expires.
use std::time::Duration;
ctx.approval(
"prod-gate",
ApprovalConfig::new("Deploy to production?")
.assigned_to(Assignee::group("release-managers"))
.with_deadline(Duration::from_secs(3600))
.on_timeout(EscalationPolicy::AutoReject),
).await?;
| Field | Builder | Meaning |
|---|---|---|
message | ApprovalConfig::new | Prompt shown to reviewers |
assignee | assigned_to | Assignee::user / Assignee::group expected to answer |
deadline_secs | with_deadline / with_deadline_secs | SLA window, in seconds |
on_timeout | on_timeout | EscalationPolicy applied when the deadline fires (defaults to AutoReject) |
timeout_seconds | with_timeout_seconds | Legacy spelling of a deadline with an implicit AutoReject |
approvers | requiring | Approvers computed by the handler: required approvals, allowed groups and an audit reason |
See Approval Gates for the full list of escalation policies, where the remaining time surfaces, and how to require several approvers.
Human input steps
#[derive(Deserialize, JsonSchema)]
struct Answers {
answers: Vec<String>,
}
let answers: Answers = ctx
.human_input("clarify", HumanInputConfig::new("Answer the clarification questions"))
.await?;
The run suspends on AwaitingApproval until a person posts an answer matching
the JSON schema of Answers to POST /api/v1/runs/:id/steps/:step_id/input.
The config takes the same deadline, escalation, assignee and approvers options
as an approval gate. A rejected input reaches the handler as
EngineError::HumanInputRejected. See Human Input.
Decision steps
A decision step asks a DecisionProvider (System One / Jev) the
questions declared by a struct and returns that struct, filled with the answers.
#[derive(DecisionChoice)]
enum Team { Billing, Technical }
#[derive(DecisionAnswers)]
struct Routing {
#[choice("Which team?")]
team: Team,
}
let routing = ctx.decision(
"triage",
DecisionConfig::new("Payouts have been failing for 3 days")
.answers::<Routing>()
.escalate_below(0.7),
).await?;
if let Team::Billing = routing.team { /* .. */ }
Below the escalate_below confidence threshold, the run suspends for human approval;
on resume the stored answers are replayed as-is. See Decisions.
Conditions
A handler branches with plain Rust if/else. That is invisible to the
execution planner, which is why two helpers
exist to declare a branch explicitly.
ctx.when(label, predicate) deserializes the run input into the type the
closure takes and evaluates the predicate on it. It returns the predicate’s
value, and the planner records it as evaluated together with the label you
gave the branch. The label is a name for the operator, never parsed:
#[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?;
}
A payload that does not match the type ("env": "prd") fails with
EngineError::Serialization instead of silently taking the else branch.
ctx.when_dynamic(label, value) declares a branch whose value comes from
a previous step’s output. It returns value unchanged; the planner records the
condition as unevaluable, because step outputs are synthetic while planning:
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?;
}
Both helpers are optional: a plain if still runs exactly the same way. They
only make the branch legible to whoever reads the plan.
Step status lifecycle
Steps follow this state machine:
stateDiagram-v2
[*] --> Pending
Pending --> Running
Running --> Completed
Running --> Failed
Pending --> Skipped
Every step transition is recorded. Failed steps report their error in the step output.
Operations
An Operation is the extensibility mechanism for custom step types. Operations let you integrate external services (GitLab, Slack, any HTTP API) as tracked steps.
The trait
use std::future::Future;
use std::pin::Pin;
use ironflow_engine::error::EngineError;
use serde_json::Value;
pub trait Operation: Send + Sync {
fn kind(&self) -> &str;
fn execute(&self) -> Pin<Box<dyn Future<Output = Result<Value, EngineError>> + Send + '_>>;
fn input(&self) -> Option<Value> { None }
}
kind()returns a short identifier (e.g."slack","gitlab") stored in the databaseexecute()runs the operation and returns JSON outputinput()optionally returns structured input for observability
Using an operation in a workflow
Operations are invoked via ctx.operation(), which takes a step name and a reference to the operation:
use ironflow_engine::context::WorkflowContext;
let slack = SlackNotify::new(&webhook_url);
ctx.operation("notify-team", &slack).await?;
Implementing an operation
use std::future::Future;
use std::pin::Pin;
use ironflow_engine::error::EngineError;
use ironflow_engine::operation::Operation;
use serde_json::{Value, json};
pub struct SlackNotify {
webhook_url: String,
message: String,
}
impl Operation for SlackNotify {
fn kind(&self) -> &str {
"slack-notify"
}
fn input(&self) -> Option<Value> {
Some(json!({ "message": self.message }))
}
fn execute(&self) -> Pin<Box<dyn Future<Output = Result<Value, EngineError>> + Send + '_>> {
Box::pin(async move {
// Send to Slack webhook using self.webhook_url
Ok(json!({ "ok": true }))
})
}
}
Built-in vs custom
Built-in step types (Shell, Http, Agent, Approval) have dedicated methods on WorkflowContext. Operations are for everything else – they give you a typed extension point without modifying the engine.
Pre-built ops crates
Ironflow ships with 13 ready-to-use ops crates under ops/ for common services: GitLab, Slack, Kubernetes, Docker, Helm, PostgreSQL, S3, Grafana, Loki, Mimir, Tempo, Git, and shared helpers. Each provides typed operations that plug directly into ctx.operation().
See Using Pre-built Ops Crates for the full catalog and usage examples, or Writing an Operation to implement your own from scratch.
Engine & Worker
Engine
The Engine is the in-memory registry that maps workflow names to handlers and orchestrates run execution. It holds references to the Store (persistence), the Provider (agent backends), and the event publisher.
let mut engine = Engine::new(store, provider);
engine.register(Box::new(MyWorkflow))?;
The Engine is used by both the API server (for metadata and describe endpoints) and the Worker (for execution).
Before running an agent step, the engine stamps the ironflow.io/run-id,
ironflow.io/root-run-id and ironflow.io/step pod labels on its config (the
step name is sanitized into a valid label value), so the Kubernetes providers
can tag the pod and clean up a previous attempt of the same step on retry. The
root run is the run itself, or the top-level run inside a sub-workflow
(ctx.root_run_id()).
Before every execution of a run (Engine::execute_handler_run, the first one
included, and Engine::resume_run after a gate under ExecutionMode::Local),
the engine calls AgentProvider::release_run with the run id. The
default does nothing; K8sEphemeralProvider deletes the pods left by a dead
attempt of the run or of its sub-workflows, and waits until they are gone. A
failed release fails the execution with a replayable error: the run goes to
Retrying while it has retries left.
Execution mode
A run suspended on an approval, a human input or an escalation resumes once the
gate is resolved. Engine::with_execution_mode decides where that happens:
ExecutionMode::Local(default): the API process moves the run toRunningand callsEngine::resume_runitself. Use it for single-process deployments where the API also registers the handlers.TestEnginealways resumes this way.ExecutionMode::Workers: the API moves the run back toPending. A worker claims it throughpick_next_pendingand finishes it withEngine::execute_handler_run, replaying the steps that already completed. Use it when the API runs without the workspace, tools or handlers the workflow needs.
let engine = Engine::new(store, provider)
.with_execution_mode(ExecutionMode::Workers);
Worker
A Worker is a background process that:
- Waits for a free execution slot (
concurrency) - Polls the API for a pending run and acquires a lease on it
- Executes the workflow handler via the Engine
- Refreshes the lease periodically during execution
- Reports the result back to the API
A saturated worker does not poll: a run is only claimed once a slot can execute it, so its lease never expires while it waits.
let worker = WorkerBuilder::new(&api_url, &worker_token)
.provider(Arc::new(ClaudeCodeProvider::new()))
.concurrency(2)
.poll_interval(Duration::from_secs(2))
.register(Box::new(MyWorkflow))
.build()?;
worker.run().await?;
Workflows that use ctx.decision(...) need a decision provider on the worker too:
.decision_provider(Arc::new(TypeSafeProvider::new(api_key))). See
Decisions.
When Provider Accounts exist, the worker picks one for
every agent step and injects its credential; WorkerBuilder::account_strategy
chooses how.
Lease & Reaper
Workers hold a time-limited lease on each run they execute. If a worker crashes or is evicted, the lease expires and the Reaper (a background task in the API server) detects the orphaned run and requeues it.
Waker
Runs paused in Sleeping (a ctx.delay step, or a ctx.wait_for_signal step waiting for its signal) carry their wake-up time in scheduled_at. The Waker, a background task of the API server, claims every due run every 10 seconds and moves it back to Pending exactly once, even with several API instances. Under ExecutionMode::Local the API then resumes the run in-process; under ExecutionMode::Workers a worker picks it up. A delivered signal wakes its runs right away, without waiting for the Waker.
Scaling
Workers are stateless. Add more workers to increase throughput. Each worker polls independently – no coordination is needed beyond the API server.
Provider Accounts
A Provider Account is an account at an AI provider: in v1, a Claude Pro/Max
subscription. Each account has a credential and usage limits (windows such as
the 5 hour and 7 day windows of a subscription). Admins manage accounts live
from the dashboard (Settings > Accounts), the REST API
(/api/v1/provider-accounts), the CLI (ironflow accounts) and the MCP server.
When at least one account of the right kind exists, the worker picks one for every agent step, injects its credential into the Claude CLI process, and records the usage windows the CLI reports. With no account, agent steps run with the worker’s own environment, exactly as before.
The claude_subscription kind
Run claude setup-token on a machine logged into the subscription and paste
the sk-ant-oat01-... token. Before storing anything, the server checks the
format and then sends a one-token request to the Anthropic API:
- a malformed or rejected token is refused with
422, and nothing is stored; - a rate-limited token (
429) is stored and shown as limited until its window resets; - an unreachable provider gives
502.
Where the credential lives
The credential is stored as the system secret accounts/<id>/credential,
encrypted like every other secret. That namespace is hidden from the Secrets
page and refused by the Secrets API. No response, log, audit entry or event
carries the token.
Injection per transport
| Transport | How the token reaches the CLI |
|---|---|
Local (ClaudeCodeProvider) | CLAUDE_CODE_OAUTH_TOKEN in the child process environment |
Docker (DockerProvider) | CLAUDE_CODE_OAUTH_TOKEN in the exec environment |
SSH (SshProvider) | first line of stdin, read by the remote shell and exported |
| Kubernetes | not yet: agent steps use the pod environment |
The token never appears on a command line. The worker forces the CLI into
stream-json mode so it can read the rate_limit_event lines that report the
windows.
Selection
The worker keeps only available accounts: enabled, token not rejected, not
expired, no applicable window rejected until its reset, and under
max_concurrency. A window scoped to a model family (for example the Opus
7 day window) only blocks steps using that family. A strategy then picks one:
| Strategy | Picks |
|---|---|
least_utilized (default) | the lowest peak utilization, plus 0.15 per running step |
priority | the lowest priority value |
round_robin | the next account, by name |
Choose it on the worker:
use ironflow_core::account_strategy::Priority;
let worker = WorkerBuilder::new(&api_url, &worker_token)
.provider(Arc::new(ClaudeCodeProvider::new()))
.account_strategy(Arc::new(Priority))
.build()?;
When every account is limited, the step fails with the time of the next reset.
CLI
claude setup-token | ironflow accounts add perso-max --token-stdin --tag perso --priority 10
ironflow accounts list
ironflow accounts usage perso-max
ironflow accounts update perso-max --max-concurrency 2
ironflow accounts test perso-max
ironflow accounts remove perso-max --yes
Usage history is kept PROVIDER_ACCOUNT_USAGE_RETENTION_DAYS days (30 by default).
Approval Gates
An approval gate pauses a workflow run until a human approves or rejects it. This enables human-in-the-loop workflows like deploy pipelines where production deploys require sign-off.
How it works
- The handler calls
ctx.approval()with a prompt message - The run transitions to
AwaitingApproval - The worker releases the run and moves on to other work
- A human calls
POST /api/v1/runs/:id/approveorPOST /api/v1/runs/:id/reject - On approval, the run resumes. Under
ExecutionMode::Workersit is requeued toPending: a worker picks it up, replays completed steps from cache, skips the approved gate, and continues execution. UnderExecutionMode::Local(the default) the API process resumes it the same way itself - On rejection, the run transitions to
Failed
If the gate carries an SLA deadline and nobody answers in time, step 4 is performed by the server instead of a human – see SLA timers below.
Example
//! Deploy workflow with human approval gate before production.
use std::time::Duration;
use ironflow_engine::config::{
ApprovalConfig, Assignee, EscalationPolicy, NotificationTarget, ShellConfig,
};
use ironflow_engine::context::WorkflowContext;
use ironflow_engine::handler::{HandlerFuture, WorkflowHandler};
/// Deploy pipeline that requires human approval before shipping to production.
///
/// 1. **build** -- compile the project
/// 2. **test** -- run the test suite
/// 3. **deploy-staging** -- deploy to staging environment
/// 4. **approval gate** -- pause and wait for human approval
/// 5. **deploy-production** -- resumes after approval via step replay
///
/// After approval (POST /api/v1/runs/:id/approve), the engine
/// re-executes the handler: completed steps return cached output,
/// the approved gate is skipped, and execution continues with
/// deploy-production.
///
/// If the approval is rejected, the run transitions to `Failed`.
/// If cancelled, the run transitions to `Cancelled`.
///
/// The gate carries a one-hour SLA. If nobody answers in time, the escalation
/// chain runs one policy per expiry: the first hour pings a webhook and keeps
/// waiting, the second gives up and fails the run. The deadline lives in the
/// database, so it survives an API or worker restart.
pub struct DeployApproval;
impl WorkflowHandler for DeployApproval {
fn name(&self) -> &str {
"deploy-approval"
}
fn description(&self) -> &str {
"Deploy pipeline with human approval gate before production. \
Demonstrates ctx.approval() for human-in-the-loop workflows."
}
fn source_code(&self) -> Option<&str> {
Some(include_str!("deploy_approval.rs"))
}
fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
Box::pin(async move {
// Step 1: Build
ctx.shell(
"build",
ShellConfig::new("echo 'Compiling...' && sleep 0.2 && echo 'Build OK'"),
)
.await?;
// Step 2: Test
ctx.shell(
"test",
ShellConfig::new("echo 'Running tests...' && sleep 0.3 && echo '87 tests passed'"),
)
.await?;
// Step 3: Deploy to staging
ctx.shell(
"deploy-staging",
ShellConfig::new(
"echo 'Deploying to staging...' && sleep 0.2 && echo 'Staging live'",
),
)
.await?;
// Step 4: Human approval gate
// The run pauses here and transitions to AwaitingApproval.
// A human must call POST /api/v1/runs/:id/approve to continue,
// or POST /api/v1/runs/:id/reject to fail the run.
ctx.approval(
"prod-approval",
ApprovalConfig::new("Staging looks good. Deploy to production?")
.assigned_to(Assignee::group("release-managers"))
.with_deadline(Duration::from_secs(3600))
.on_timeout(EscalationPolicy::Chain(vec![
// After 1 h without an answer: warn, keep waiting.
EscalationPolicy::Notify(vec![NotificationTarget::Webhook {
url: "https://example.com/hooks/deploy-sla".to_string(),
}]),
// After 2 h: give up rather than ship unreviewed.
EscalationPolicy::AutoReject,
])),
)
.await?;
// Step 5: Deploy to production (only reached after approval)
ctx.shell(
"deploy-production",
ShellConfig::new(
"echo 'Deploying to production...' && sleep 0.3 && echo 'Production live'",
),
)
.await?;
Ok(())
})
}
}
Configuration
use std::time::Duration;
ApprovalConfig::new("Deploy to production?")
.assigned_to(Assignee::group("release-managers")) // Who is expected to answer
.with_deadline(Duration::from_secs(3600)) // SLA: one hour to answer
.on_timeout(EscalationPolicy::AutoReject) // What happens when it expires
assigned_to takes an Assignee – Assignee::user("alice") or
Assignee::group("release-managers"). It drives notification routing and audit,
and it decides who may resolve the gate:
- an admin resolves any gate;
- a gate assigned to a user is resolved by that user, admin or not, and by whoever holds an active delegation from them;
- a gate assigned to a group, or to nobody, is admin-only.
The assignee is matched to the caller by user ID, so an API key resolves its owner’s gates whatever the key is named.
Everything past the message is optional. Without a deadline, the run waits indefinitely.
SLA timers
with_deadline (or with_deadline_secs) arms a timer on the approval step. The
deadline is stored in the database next to the step, not in memory, so it
survives an API or worker restart: a fresh process picks the expired gate up on
its next pass. The API server checks for expired gates every 30 seconds.
The timer is cleared the moment the gate resolves – approved, rejected, or escalated – so a gate is never escalated after a human answered it. A deadline fires at most once, even with several API instances running.
with_timeout_seconds is the legacy spelling: it is now enforced, as a
deadline with an implicit AutoReject policy. Setting both keeps the explicit
with_deadline.
Escalation policies
on_timeout takes an EscalationPolicy. Without one, an expired deadline
auto-rejects.
| Policy | What it does when the deadline fires |
|---|---|
AutoApprove | Completes the gate with approved_by: "system:timeout" and resumes the run: in the server under ExecutionMode::Local, by requeuing it for a worker under ExecutionMode::Workers (execution mode). |
AutoReject | Fails the step and the run with approval timeout. The default. |
Notify(targets) | Posts the escalation event to each target, leaves the gate open, restarts the timer. |
Escalate(Assignee) | Reassigns the gate to another user or group, leaves it open, restarts the timer. |
Chain(policies) | Applies one policy per expiry, in order. |
Notify and Escalate do not resolve the gate: on their own, they fire again at
every expiry until a human answers. Each firing writes an audit entry, so the
loop is visible rather than silent. Wrap them in a Chain to advance one policy
per expiry instead:
use std::time::Duration;
ApprovalConfig::new("Deploy to production?")
.with_deadline(Duration::from_secs(3600))
.on_timeout(EscalationPolicy::Chain(vec![
// After 1 h: ping the on-call channel, keep waiting.
EscalationPolicy::Notify(vec![NotificationTarget::Slack {
webhook_url: slack_webhook_url,
channel: "#deploys".to_string(),
}]),
// After 2 h: give up.
EscalationPolicy::AutoReject,
]))
Once a chain runs out, the gate stays open with no timer and a warning is logged – it is never silently auto-rejected.
NotificationTarget is delivered as a plain HTTP POST, with the engine’s
shared retry and backoff: Webhook { url } posts the escalation event as JSON,
Slack { webhook_url, channel } posts a message to a Slack incoming webhook.
A dead endpoint is logged and never blocks the timer reset.
Seeing the remaining time
The countdown surfaces in three places:
- the API:
approval_seconds_remainingandapproval_assigneeon every step ofGET /api/v1/runs/:id(clamped at 0,nullwithout a deadline); - the dashboard: a countdown badge on the gate in the run’s step list;
- the CLI: the
SLAcolumn ofironflow run steps <id>, yellow in the last tenth of the window and red once expired.
Every escalation is also recorded in the audit log as an approval_escalated
event carrying the stage, the policy, what it did, and why it fired.
Delegation and absence
An approval gate assigned to one person stops every run behind it the moment that person is away. A delegation hands their approval power to a colleague for a bounded window, without making anyone an admin and without reassigning the gates one by one.
A delegation records who grants it, who receives it, the window it is valid for, and an optional glob on the workflow name:
# Alice hands her deploy approvals to Bob for a week.
ironflow-cli delegation create <bob-user-id> \
--until 2026-10-01T00:00:00Z \
--workflow 'deploy-*'
# Everything Alice granted, plus everything she received (20 per page).
ironflow-cli delegation list --page 1 --per-page 20
# Back early.
ironflow-cli delegation delete <delegation-id>
The same three endpoints back the CLI:
| Endpoint | What it does |
|---|---|
POST /api/v1/approval-delegations | Grant a delegation. The delegator is always the caller. |
GET /api/v1/approval-delegations | List the active delegations you granted or received, paginated with page and per_page (default 20, max 100). An admin sees them all and may filter with from_user_id and to_user_id. |
DELETE /api/v1/approval-delegations/{id} | Revoke one. Only the delegator or an admin may. |
What a delegation covers
Only a gate assigned to an individual can be delegated:
ApprovalConfig::new("Deploy to production?")
.assigned_to(Assignee::user("alice")) // Bob can answer this through a delegation.
ApprovalConfig::new("Deploy to production?")
.assigned_to(Assignee::group("release-managers")) // Admin-only; no single delegator.
A gate assigned to a group, or with no assignee at all, stays admin-only: there is no single person whose power could have been handed over.
The workflow_filter glob narrows a delegation to part of the catalogue.
"deploy-*" covers deploy-prod but not cleanup; omitting it covers every
workflow. A pattern that does not parse matches nothing, so a corrupted row can
never widen someone’s reach.
Several delegations can be active at once – one per colleague, one per workflow family, or from several delegators to the same person. The newest one that matches both the gate’s assignee and the run’s workflow wins.
Expiry
The window is half-open: a delegation is live at valid_from and already over
at valid_until. There is no cleanup job. Expired and not-yet-started rows are
filtered out every time delegations are read, so they can neither be listed nor
used to approve – they are still reachable by ID, which is what makes an
expired delegation revocable.
Audit
A delegated decision names both people. The approval_granted (or
approval_rejected) audit entry reads:
{
"type": "approval_granted",
"run_id": "01932f...",
"approved_by": "bob (delegated from alice)",
"at": "2026-09-22T10:15:00Z"
}
An admin, or the assignee resolving their own gate, is recorded under their own name alone.
Requiring several approvers
A gate can require more than one approval, and restrict who may vote, depending
on the run itself: a small payment needs one approver, a large one needs two
people from finance. The handler decides in plain Rust, from its typed input and
the outputs of earlier steps, and passes the result to requiring:
use ironflow_engine::config::{ApprovalConfig, Approvers};
let payment: Payment = ctx.input().await?;
let approvers = match payment.amount {
a if a > 100_000 => Approvers::at_least(3)
.from_groups(["finance", "board"])
.because("amount > 100k"),
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?;
| Builder | Meaning |
|---|---|
Approvers::any() | One approval, from anyone allowed to answer the gate |
Approvers::at_least(n) | n distinct approvals (required_approvers) |
.from_groups([..]) | Only members of these groups may vote (approver_groups) |
.because("..") | Audit label shown on the dashboard and in the events (reason), never evaluated |
A typo in a field name or a comparison does not compile, and the compiler checks
every branch of the match. A gate without requiring behaves as a single
approval gate.
The approvers are stored on the step as an ApprovalRequirement (reason,
required_approvers, approver_groups) when the gate opens. That record is the
source of truth from then on: replaying or resuming the run never recomputes
it, even if the handler would now compute other approvers. GET /api/v1/runs/:id exposes it on the step as approval_requirement, with the
votes cast so far in approvals and the count needed in approvals_required;
the dashboard shows it as an n/m approvals badge whose tooltip gives the
reason.
Approvers::at_least(0) and a blank group name panic, so a broken gate fails
when the workflow runs the builder, not when a human votes. A JSON config with
required_approvers: 0 is rejected on deserialization.
Voting
- One vote per user. Votes are counted by user ID: an API key votes as its
owner, and the same user approving twice gets
409 Conflict. - An admin’s approval is one vote. Admins may always vote, even on a gate restricted to groups, but they do not override the count.
- A rejection vetoes. One rejection from anyone allowed to vote fails the run, even after partial approvals.
- Until the count is reached,
POST /approvereturns200with the run stillawaiting_approval, the gate keeps its SLA timer, and the CLI printsApproval recorded; more approvals are required.
An EscalationPolicy::AutoApprove still resolves the gate outright, whatever
the required count.
Approver groups
When the approvers list approver_groups, only members of at least one of
those groups (and admins) may vote. The gate’s assignee and approval
delegations are not consulted. A listed group without members leaves the gate
to admins.
Group membership is managed by admins:
# Put alice in finance and legal (replaces her current groups).
ironflow user set-groups <alice-id> --group finance --group legal
# Show her groups.
ironflow user groups <alice-id>
# Remove her from every group.
ironflow user set-groups <alice-id>
The same operations are available as GET and PUT /api/v1/users/:id/groups.
Group names are 1 to 64 characters from [A-Za-z0-9_.-], at most 50 per user.
Audit events
approval_requestedis published when the gate opens and carries the recordedrequirement.approval_grantedis published for every vote, with thestep_id,approvals_received,approvals_requiredand therequirement. The gate resolves whenapprovals_received >= approvals_required.approval_rejectedcarries thestep_idand therequirement.
{
"type": "approval_granted",
"run_id": "01932f...",
"step_id": "01932f...",
"approved_by": "alice",
"approvals_received": 1,
"approvals_required": 2,
"requirement": {
"reason": "amount > 10k",
"required_approvers": 2,
"approver_groups": ["finance"]
},
"at": "2026-09-24T10:15:00Z"
}
Step replay
After an approval, the engine re-executes the handler from the beginning: in the API process under ExecutionMode::Local, on the worker that picks up the requeued run under ExecutionMode::Workers (see execution mode). Completed steps return their cached output immediately – they do not re-run. The approved gate is skipped, and execution resumes with the next step.
Human Input
A human input step pauses a workflow run until a person submits a typed answer. Where an approval gate asks “yes or no?”, a human input asks for data: answers to clarification questions, a choice, a value. The handler gets the answer back as a Rust type.
How it works
- The handler calls
ctx.human_input::<T>()with a message.TderivesDeserializeandJsonSchema. - The engine records a step of kind
human_inputwhose input holds the message and the JSON schema ofT(under the keyschema). - The step and the run move to
AwaitingApproval, the same status as an approval gate, and aninput_requiredevent is published on the run’s event stream. - A person posts an answer to
POST /api/v1/runs/:id/steps/:step_id/input. The API validates it against the stored schema, completes the step and resumes the run: in the API process underExecutionMode::Local(the default), or by requeuing it toPendingfor a worker underExecutionMode::Workers(see execution mode). - The handler is replayed: completed steps come from cache and
human_inputreturns the answer, deserialized intoT.
Example
use ironflow_engine::prelude::*;
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?;
ctx.shell("plan", ShellConfig::new(&format!("./plan.sh {}", answers.answers.len())))
.await?;
Configuration
HumanInputConfig reuses the approval gate machinery:
use std::time::Duration;
HumanInputConfig::new("Which environment should we target?")
.assigned_to(Assignee::user("alice")) // Who is expected to answer
.requiring(Approvers::any().from_groups(["product"])) // Who may answer
.with_deadline(Duration::from_secs(3600)) // SLA: one hour to answer
.on_timeout(EscalationPolicy::AutoReject) // What happens when it expires
| Field | Builder | Meaning |
|---|---|---|
message | HumanInputConfig::new | Prompt shown to the person answering |
assignee | assigned_to | Assignee::user / Assignee::group expected to answer |
approvers | requiring | Groups allowed to answer. The first valid answer wins, whatever the count |
deadline_secs | with_deadline / with_deadline_secs | SLA window, in seconds |
on_timeout | on_timeout | EscalationPolicy applied when the deadline fires (defaults to AutoReject) |
Who may answer follows the approval rules: an admin, a member of the requiring
groups, the assignee, or someone holding a delegation from the assignee. The
person who answered is recorded on the step like an approval vote.
EscalationPolicy::AutoApprove has no meaning without a value:
on_timeout(EscalationPolicy::AutoApprove) panics, and so does a Chain
containing it. Notify, Escalate, AutoReject and Chain behave as for
approval gates.
API
Answer the input with a body matching the schema:
curl -X POST "$IRONFLOW_URL/api/v1/runs/$RUN_ID/steps/$STEP_ID/input" \
-H "Authorization: Bearer $TOKEN" \
-H "Content-Type: application/json" \
-d '{"answers": ["staging", "eu-west-1"]}'
-
200returns the run, nowrunning. -
422with codeINVALID_INPUTwhen the body does not match the schema; each violation is listed inerror.details.errors:{ "error": { "code": "INVALID_INPUT", "message": "input does not match the expected schema", "details": { "errors": ["3 is not of type \"array\""] } } } -
409when the input was already answered or rejected. -
400when the step is not a human input, or the run is not waiting on it.
Refuse the input, with an optional reason:
curl -X POST "$IRONFLOW_URL/api/v1/runs/$RUN_ID/steps/$STEP_ID/reject" \
-H "Authorization: Bearer $TOKEN" \
-H "Content-Type: application/json" \
-d '{"reason": "out of scope for this sprint"}'
POST /api/v1/runs/:id/approve refuses a run waiting on a human input with a
400: approving it would resume the handler without an answer.
The dashboard shows a form for every pending input on the run page. The CLI has
ironflow run input <run> <step> --value '{..}' (or --value-file) and
ironflow run reject-input <run> <step> --reason ..; the MCP server has the
submit_input and reject_input tools.
Rejection
A rejected input does not fail the run by itself. The step is marked Rejected
with the reason, the run resumes, and human_input returns
EngineError::HumanInputRejected, so the handler decides what happens next:
match ctx.human_input::<Answers>("clarify", config).await {
Ok(answers) => { /* use the answers */ }
Err(EngineError::HumanInputRejected { reason, .. }) => {
ctx.shell("notify", ShellConfig::new(&format!("./notify.sh '{reason}'")))
.await?;
}
Err(err) => return Err(err),
}
A handler that propagates the error fails the run; it is never retried
automatically. The run-level POST /api/v1/runs/:id/reject still fails the run
outright.
Replay and retries
- A run resumed without an answer suspends again on the same step; no new step is created.
- An answer given before an automatic retry is carried over to the next attempt: the person is not asked twice.
- If the handler changed and the stored answer no longer fits
T, the run fails with a step configuration error.
Events
input_required is published on GET /api/v1/runs/:id/events when the step
opens. It carries run_id, step_id, step_name, step_index, message and
schema, everything a client needs to render a form.
Execution plans
Planning never suspends. The step is recorded with kind human_input, and T
is built from {}: a type with #[serde(default)] lets the plan continue past
the input. Otherwise the plan stops there with the reason
human input '<name>' has no answer while planning.
Testing
TestEngine::with_mock_human_input answers every human input without waiting:
use ironflow_engine::testing::{HumanInputOutcome, TestEngine};
use serde_json::json;
let result = TestEngine::new()
.with_handler(Clarify)
.with_mock_human_input(|_name, _config| {
HumanInputOutcome::Provided(json!({"answers": ["staging"]}))
})
.run(json!({}))
.await?;
HumanInputOutcome::reject("reason") makes the handler receive
EngineError::HumanInputRejected. Without the mock, the run ends in
AwaitingApproval: write the answer on the step through the store, then call
TestEngine::resume.
Signals
A signal is an external message sent to Ironflow to resume the runs waiting
for it. It has a name (what happened, e.g. ci.pipeline_finished) and a
key (which occurrence, e.g. a commit SHA). A run waits with
ctx.wait_for_signal; a producer delivers with POST /api/v1/signals, the CLI,
the MCP server or Engine::send_signal.
Declaring a signal
A signal is a typed payload: a struct implementing Signal, which names it.
use ironflow_engine::signal::Signal;
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
#[derive(Serialize, Deserialize, JsonSchema)]
struct PipelineFinished {
status: String,
}
impl Signal for PipelineFinished {
const NAME: &'static str = "ci.pipeline_finished";
}
Waiting for a signal
let finished = ctx
.wait_for_signal::<PipelineFinished>("wait-ci", &sha, Duration::from_secs(3600))
.await?;
match finished {
Some(pipeline) if pipeline.status == "success" => { /* deploy */ }
Some(_) => return Err(EngineError::StepConfig("CI failed".to_string())),
None => return Err(EngineError::StepConfig("CI timed out".to_string())),
}
- The key is the occurrence: wait on a commit SHA, not on a merge request. A signal for an older push must never resume a run waiting for the newest one.
- Received early: a signal delivered after the run was created but before the step opened resolves the step at once, without suspending.
- Suspended: otherwise the run goes
Sleepinguntil the deadline. A delivery whose payload matches the JSON schema of the type wakes it, and the handler receivesSome(payload). - Timeout: when the deadline passes first, the step completes as timed out and
the handler receives
None. - Broadcast: every run waiting on the same
(name, key)receives the signal. - Invalid payload: a run whose schema the payload does not match keeps waiting
and is listed under
rejectedin the response. The signal is still stored. - Idempotency: a second delivery with the same
idempotency_idreturnsduplicate: trueand delivers nothing.
Woken runs resume in-process under ExecutionMode::Local, or go back to
Pending for a worker under ExecutionMode::Workers. Timeouts are applied by the
waker task of the API server (see Engine and worker).
Sending a signal
| Channel | How |
|---|---|
| REST | POST /api/v1/signals with {"name", "key", "payload", "idempotency_id"}; admin JWT or API key with the signals_send scope |
| REST | GET /api/v1/signals?name=&key= lists received signals (runs_read for an API key) |
| CLI | ironflow signal send ci.pipeline_finished --key <sha> --payload '{"status":"success"}' |
| CLI | ironflow signal list --name ci.pipeline_finished |
| MCP | send_signal, list_signals |
| Rust | engine.send_signal(&PipelineFinished { .. }, &sha, Some(&delivery_id)) |
Example: wait for CI
The workflow pushes a commit and waits for its pipeline:
ctx.shell("push", ShellConfig::new("git push origin HEAD")).await?;
let finished = ctx
.wait_for_signal::<PipelineFinished>("wait-ci", &sha, Duration::from_secs(3600))
.await?;
The CI webhook handler filters the terminal pipeline statuses and forwards them, using the delivery ID of the webhook so a retried delivery is not counted twice:
if matches!(status.as_str(), "success" | "failed" | "canceled") {
engine
.send_signal(&PipelineFinished { status }, &sha, Some(&delivery_id))
.await?;
}
Signals are kept SIGNAL_RETENTION_DAYS days (default 7) and then purged. The example server wires the variable into its RunPurger; a custom server builds it with RunPurger::from_config(store, &config).
Decisions
A decision is a typed machine verdict: classify, route, score, or answer yes/no, with a calibrated confidence instead of free text. It is the shape of TypeSafe AI’s System One model (Jev): for a structured verdict it is far faster and cheaper than routing a prompt through a conversational LLM.
Decisions use a dedicated abstraction rather than the agent one: there is
no single prompt, no tools, and no streaming. A run wires a DecisionProvider
independently of its agent provider.
Wiring a provider
use ironflow_core::providers::http::TypeSafeProvider; // feature "provider-typesafe"
use std::sync::Arc;
let engine = Engine::new(store, agent_provider)
.with_decision_provider(Arc::new(TypeSafeProvider::new(api_key)));
When the workflow runs in a worker, wire the provider on the
WorkerBuilder instead; every run the worker executes gets it:
let worker = WorkerBuilder::new(&api_url, &worker_token)
.provider(agent_provider)
.decision_provider(Arc::new(TypeSafeProvider::new(api_key)))
.build()?;
Without a provider, a decision step fails with NoDecisionProvider. In tests, use
RecordReplayDecisionProvider::replay(dir) to serve captured JSON fixtures with no
network.
Through OpenRouter
OpenRouter serves the same System One wire contract on its Decisions endpoint, so
the same provider reaches Jev with only the base URL and key changing. Use the
openrouter constructor and select the OpenRouter model slug on the config
(typesafe/jev-1.13, exposed as OPENROUTER_MODEL). OpenRouter requires this
concrete versioned slug; the jev-latest alias returns 400 "Model does not exist":
use ironflow_core::providers::http::typesafe::OPENROUTER_MODEL;
use ironflow_core::providers::http::TypeSafeProvider;
use ironflow_engine::config::DecisionConfig;
use std::sync::Arc;
let engine = Engine::new(store, agent_provider)
.with_decision_provider(Arc::new(TypeSafeProvider::openrouter(openrouter_key)));
let config = DecisionConfig::new(state).model(OPENROUTER_MODEL);
OpenRouter’s decisions route is on an alpha path that may move; override it with
TypeSafeProvider::with_endpoint(url) if it relocates. OpenRouter also requires
question instructions and criteria to be strings, which the derive’s string
literals already satisfy.
The three question types
The questions are the fields of a struct deriving DecisionAnswers; the options of a
choice are the unit variants of an enum deriving DecisionChoice. ctx.decision
returns the struct itself.
| Type | Field attribute | Field type | Answer |
|---|---|---|---|
| noul | #[noul("..")], optionally if_true = "..", if_false = ".." | f64 | probability of “yes” in [0, 1] |
| choice | #[choice("..")] | an enum deriving DecisionChoice | the option picked |
| score | #[score("..", levels = ["..", ".."])] | f64 | probability-weighted level index |
use ironflow_engine::config::DecisionConfig;
use ironflow_engine::decision::{DecisionAnswers, DecisionChoice};
#[derive(DecisionChoice)]
enum Team {
#[choice(description = "Payments, invoices, refunds")]
Billing,
Technical,
Sales,
}
#[derive(DecisionAnswers)]
struct Triage {
#[noul("Does this convey urgency?")]
is_urgent: f64,
#[choice("Which team?")]
team: Team,
#[score("How frustrated?", levels = ["Calm", "Frustrated", "Very angry"])]
mood: f64,
}
let triage = ctx.decision(
"triage",
DecisionConfig::new("Payouts have been failing for 3 days")
.answers::<Triage>()
.escalate_below(0.7),
).await?;
match triage.team {
Team::Billing => { /* .. */ }
Team::Technical | Team::Sales => { /* .. */ }
}
A question is named after its field. An option is labelled with its variant name in
snake_case; #[choice(rename = "..")] changes the label and
#[choice(description = "..")] tells the model what the option means (doc comments are
never sent). A field without a question, a score without levels or a field of the wrong
type does not compile. An option the provider returns that is not a variant fails the
step with DecisionError::UnknownChoice.
The options are fixed at compile time. choice and score answers carry a
confidence; a noul answer reports only a probability p, and its confidence is
derived as 2 * |p - 0.5| (a coin flip is 0, a certain yes/no is 1). Confidence
drives escalation, below; it is not part of the typed answer.
Escalation
When escalate_below(threshold) is set and any answer’s confidence falls below it,
the run suspends in AwaitingApproval, exactly like an approval gate.
On resume the decision is not re-run: the stored answers are replayed as-is, so
downstream routing stays deterministic across the suspend/resume boundary.
Cost
The provider reports token usage; the engine imputes the cost to the run’s budget like an agent step. For Jev, only input tokens are billed (output is unmetered).
Writing a Workflow
This guide walks through creating a workflow handler from scratch.
1. Define your input
If your workflow accepts input, define a struct with Deserialize and JsonSchema:
use schemars::JsonSchema;
use serde::Deserialize;
#[derive(Deserialize, JsonSchema)]
struct DeployInput {
environment: String,
version: String,
}
The JsonSchema derive lets the dashboard render a dynamic form for triggering the workflow.
2. Implement WorkflowHandler
use ironflow_engine::config::ShellConfig;
use ironflow_engine::context::WorkflowContext;
use ironflow_engine::handler::{HandlerFuture, WorkflowHandler, input_schema_for};
use serde_json::Value;
pub struct Deploy;
impl WorkflowHandler for Deploy {
fn name(&self) -> &str {
"deploy"
}
fn description(&self) -> &str {
"Deploy a version to an environment"
}
fn input_schema(&self) -> Option<Value> {
Some(input_schema_for::<DeployInput>())
}
fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
Box::pin(async move {
let input: DeployInput = ctx.input().await?;
ctx.shell(
"build",
ShellConfig::new(&format!("echo 'Building {}'", input.version)),
).await?;
ctx.shell(
"deploy",
ShellConfig::new(&format!(
"echo 'Deploying {} to {}'",
input.version, input.environment
)),
).await?;
Ok(())
})
}
}
Branching is plain Rust if/else. When a branch depends on the run input,
declare it with ctx.when("production run", |i: &DeployInput| i.environment == "production")
so it shows up in the execution plan; use ctx.when_dynamic when
the branch depends on a previous step’s output. Run ironflow run plan <name> --input '{}' to see the steps your handler would create before triggering it.
3. Register in your handlers list
pub fn handlers() -> Vec<Box<dyn WorkflowHandler>> {
vec![
Box::new(Deploy),
// ... other handlers
]
}
Both the server and the worker must register the same handlers. The recommended pattern is a shared handlers() function in a library crate.
4. Complete example
The greeting workflow in the examples directory demonstrates all features:
use std::collections::HashMap;
use ironflow_engine::config::ShellConfig;
use ironflow_engine::context::WorkflowContext;
use ironflow_engine::handler::{HandlerFuture, WorkflowHandler, input_schema_for};
use schemars::JsonSchema;
use serde::Deserialize;
use serde_json::Value;
/// Input payload for the greeting workflow.
///
/// Derives [`JsonSchema`] so the dashboard can render a dynamic form.
#[derive(Deserialize, JsonSchema)]
struct GreetingInput {
/// Person to greet.
name: String,
/// Greeting language (en, fr, es).
#[serde(default = "default_language")]
language: String,
/// Number of times to repeat the greeting.
#[serde(default = "default_repeat")]
repeat: u32,
/// Whether to output in uppercase.
#[serde(default)]
uppercase: bool,
}
fn default_language() -> String {
"en".to_string()
}
fn default_repeat() -> u32 {
1
}
pub struct Greeting;
impl WorkflowHandler for Greeting {
fn name(&self) -> &str {
"greeting"
}
fn category(&self) -> Option<&str> {
Some("examples")
}
fn input_schema(&self) -> Option<Value> {
Some(input_schema_for::<GreetingInput>())
}
fn default_labels(&self) -> HashMap<String, String> {
HashMap::from([("project".to_string(), "ironflow".to_string())])
}
fn description(&self) -> &str {
"A demo workflow that greets someone. \
Shows how input_schema generates a dynamic form in the dashboard."
}
fn source_code(&self) -> Option<&str> {
Some(include_str!("greeting.rs"))
}
fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
Box::pin(async move {
let input: GreetingInput = ctx.input().await?;
let greeting = match input.language.as_str() {
"fr" => format!("Bonjour, {} !", input.name),
"es" => format!("Hola, {}!", input.name),
_ => format!("Hello, {}!", input.name),
};
let mut message = (0..input.repeat)
.map(|_| greeting.as_str())
.collect::<Vec<_>>()
.join("\n");
if input.uppercase {
message = message.to_uppercase();
}
ctx.shell("greet", ShellConfig::new(&format!("echo '{message}'")))
.await?;
Ok(())
})
}
}
Writing an Operation
Operations let you extend Ironflow with custom step types for integrating external services.
1. Implement the Operation trait
use std::env;
use std::future::Future;
use std::pin::Pin;
use ironflow_engine::error::EngineError;
use ironflow_engine::operation::Operation;
use serde_json::{Value, json};
pub struct SlackNotify {
webhook_url: String,
message: String,
}
impl SlackNotify {
pub fn new(message: &str) -> Self {
let webhook_url = env::var("SLACK_WEBHOOK_URL")
.expect("SLACK_WEBHOOK_URL env var required");
Self {
webhook_url,
message: message.to_string(),
}
}
}
impl Operation for SlackNotify {
fn kind(&self) -> &str {
"slack-notify"
}
fn input(&self) -> Option<Value> {
Some(json!({ "message": self.message }))
}
fn execute(&self) -> Pin<Box<dyn Future<Output = Result<Value, EngineError>> + Send + '_>> {
Box::pin(async move {
let client = reqwest::Client::new();
let resp = client
.post(&self.webhook_url)
.json(&json!({ "text": self.message }))
.send()
.await
.map_err(|e| EngineError::OperationFailed {
kind: "slack-notify".to_string(),
message: e.to_string(),
})?;
Ok(json!({ "status": resp.status().as_u16() }))
})
}
}
2. Use it in a workflow
let notifier = SlackNotify::new("Deploy complete!");
ctx.operation("notify-team", ¬ifier).await?;
The step is tracked in the database like any other step, with its input, output, and status.
Using a Pre-built Ops Crate
Ironflow ships with 13 ops crates under ops/ that provide ready-to-use integrations for common services. Instead of implementing the Operation trait yourself, you can use these crates to get typed, tracked operations in a single cargo add.
The pattern
Every ops crate follows the same three-step pattern:
- Add the dependency to your workflow crate
- Build a client from your workflow’s
OperationContext(viafrom_context()) - Run a tracked operation via
ctx.operation()
use ironflow_ops_slack::SlackClient;
use ironflow_ops_slack::chat::ChatPostMessage;
use slack_morphism::api::SlackApiChatPostMessageRequest;
use slack_morphism::{SlackChannelId, SlackMessageContent};
// 1. Build the client (reads slack_bot_token from the secret store)
let slack = SlackClient::from_context(&ctx).await?;
// 2. Create the operation
let req = SlackApiChatPostMessageRequest::new(
SlackChannelId::new("#deployments".to_string()),
SlackMessageContent::new().with_text("Deploy complete".to_string()),
);
let op = ChatPostMessage::new(&slack, req);
// 3. Execute as a tracked workflow step
let output = ctx.operation("notify-team", &op).await?;
The step is tracked in the database with its input, output, kind, and status, just like a shell or HTTP step.
Available crates
| Crate | Description | cargo add |
|---|---|---|
ironflow-ops-common | Shared HTTP client and helpers for ops crates (not used directly in workflows) | cargo add ironflow-ops-common |
ironflow-ops-docker | Docker containers, images, networks, volumes via bollard | cargo add ironflow-ops-docker |
ironflow-ops-git | Git operations (commit, branch, merge, diff, …) via git2 | cargo add ironflow-ops-git |
ironflow-ops-gitlab | GitLab API v4 (issues, MRs, pipelines, …) via the gitlab crate | cargo add ironflow-ops-gitlab |
ironflow-ops-grafana | Grafana API (dashboards, alerting, data sources, …) | cargo add ironflow-ops-grafana |
ironflow-ops-helm | Helm CLI wrapper (install, upgrade, rollback, charts, repos) | cargo add ironflow-ops-helm |
ironflow-ops-k8s | Kubernetes typed API via kube + k8s-openapi, plus run-to-completion ops (PodRun, JobRun, ApplyConfigMap, ApplySecret) | cargo add ironflow-ops-k8s |
ironflow-ops-loki | Grafana Loki (log queries, ingest, rules, labels) | cargo add ironflow-ops-loki |
ironflow-ops-mimir | Grafana Mimir (PromQL queries, remote write, rules, cardinality) | cargo add ironflow-ops-mimir |
ironflow-ops-postgres | PostgreSQL queries and admin via sqlx | cargo add ironflow-ops-postgres |
ironflow-ops-s3 | AWS S3 objects, buckets, presigned URLs via aws-sdk-s3 | cargo add ironflow-ops-s3 |
ironflow-ops-slack | Slack API (chat, conversations, files, users) via slack-morphism | cargo add ironflow-ops-slack |
ironflow-ops-tempo | Grafana Tempo (trace queries, search, metrics, cluster) | cargo add ironflow-ops-tempo |
Example: GitLab integration
use ironflow_ops_gitlab::GitLab;
use gitlab::api::projects::issues::CreateIssue;
// Build from workflow context (reads gitlab_token from the secret store)
let gitlab = GitLab::from_context(&ctx).await?;
// For a self-hosted instance
let gitlab = GitLab::from_context_with_host(&ctx, "gitlab.example.com").await?;
// Create an issue as a tracked step
let endpoint = CreateIssue::builder()
.project("my-group/my-project")
.title("Automated bug report")
.description("Detected by workflow")
.build()?;
let output = ctx.operation("create-issue", &gitlab.op(endpoint)).await?;
Example: Kubernetes + Helm
use ironflow_ops_k8s::{KubeClient, verb};
use ironflow_ops_helm::HelmClient;
use ironflow_ops_helm::release::Upgrade;
use k8s_openapi::api::apps::v1::Deployment;
// Check deployment status
let kube = KubeClient::from_context(&ctx).await?;
let deployments = kube.namespaced::<Deployment>("production");
let op = kube.op(deployments, verb::Get::new("my-app"));
ctx.operation("check-deployment", &op).await?;
// Upgrade the Helm release
let helm = HelmClient::from_context(&ctx).await?;
let upgrade = Upgrade::new(helm, "my-app", "charts/my-app")
.namespace("production")
.set("image.tag", "v2.1.0");
ctx.operation("upgrade-release", &upgrade).await?;
Secrets and authentication
Each crate documents its required and optional secrets in its README. The general pattern is:
- Register secrets in your workflow’s secret store
- Call
from_context(&ctx)which resolves them automatically - The client handles authentication transparently
See each crate’s README for the exact secret names and formats.
Writing your own operation
If none of the pre-built crates cover your service, see Writing an Operation to implement the Operation trait directly.
Parallel Execution
Ironflow supports running multiple steps in parallel within a workflow.
Using ctx.parallel()
Pass a list of step configurations to ctx.parallel(). All steps run concurrently and the method returns when all complete:
//! Demo workflow showcasing parallel execution and conditional branching.
use ironflow_engine::config::{ShellConfig, StepConfig};
use ironflow_engine::context::WorkflowContext;
use ironflow_engine::handler::{HandlerFuture, WorkflowHandler};
/// Simulated CI pipeline that demonstrates DAG features:
///
/// 1. **build** (sequential)
/// 2. **test-unit + test-integration + lint** (parallel)
/// 3. **deploy** or **notify-failure** (conditional branch on test results)
pub struct CiPipeline;
impl WorkflowHandler for CiPipeline {
fn name(&self) -> &str {
"ci-pipeline"
}
fn description(&self) -> &str {
"Simulated CI pipeline with parallel tests and conditional deploy. \
Demonstrates ctx.parallel() and native Rust if/else branching."
}
fn source_code(&self) -> Option<&str> {
Some(include_str!("ci_pipeline.rs"))
}
fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
Box::pin(async move {
// Step 1: Build
let build = ctx
.shell(
"build",
ShellConfig::new(
"echo 'Compiling project...' && sleep 0.2 && echo 'Build successful'",
),
)
.await?;
if !build.is_success() {
ctx.shell(
"notify-build-failure",
ShellConfig::new("echo 'BUILD FAILED - notifying team'"),
)
.await?;
return Ok(());
}
// Step 2: Parallel tests + lint
let results = ctx
.parallel(
vec![
(
"test-unit",
StepConfig::Shell(ShellConfig::new(
"echo 'Running unit tests...' && sleep 0.3 && echo '42 tests passed'",
)),
),
(
"test-integration",
StepConfig::Shell(ShellConfig::new(
"echo 'Running integration tests...' && sleep 0.5 && echo '12 tests passed'",
)),
),
(
"lint",
StepConfig::Shell(ShellConfig::new(
"echo 'Running linter...' && sleep 0.1 && echo 'No warnings'",
)),
),
],
true,
)
.await?;
// Step 3: Conditional deploy
let all_passed = results.iter().all(|r| r.output.is_success());
if all_passed {
ctx.shell(
"deploy",
ShellConfig::new("echo 'Deploying to production...' && sleep 0.2 && echo 'Deployed successfully'"),
)
.await?;
} else {
ctx.shell(
"notify-test-failure",
ShellConfig::new("echo 'TESTS FAILED - deployment skipped'"),
)
.await?;
}
Ok(())
})
}
}
How it works
- All steps in a
parallel()call start at the same time - The method returns a
Vec<ParallelStepResult>with outputs in the same order as the input; eachoutputreads through the same typed accessors as a single step (stdout(),status(),body(),is_success()) and hands out the artifacts its step declared (r.output.artifact("report.html")?) - If
fail_fastistrue(the second argument), the remaining steps are cancelled when one fails - If
fail_fastisfalse, all steps run to completion regardless of individual failures - Every step of a wave needs its own name: a wave with two steps of the same name fails
with
EngineError::StepConfigbefore anything runs (a dry-run plan reports it too) - When a run resumes (after an approval, a human input or a delay), the steps of the wave that already completed are replayed from the store; only the others run again
Conditional branching
Since workflows are Rust code, conditional logic is just if/else:
let results = ctx.parallel(steps, true).await?;
let all_passed = results.iter().all(|r| r.output.is_success());
if all_passed {
ctx.shell("deploy", ShellConfig::new("echo 'Deploying'")).await?;
} else {
ctx.shell("notify", ShellConfig::new("echo 'Tests failed'")).await?;
}
No special DSL for branching – Rust control flow works directly.
Execution Plans
An execution plan answers one question before you trigger anything: what would this workflow do with this input? It lists the steps the run would create, in order, with their dependencies, their parallel waves, their branch conditions and — when the workflow has a history — an estimate of how long each one takes.
What plan mode does, and does not
Ironflow workflows are Rust-native handlers, not declarative graphs. There is no static file to read, so the only way to know which steps a run would create is to execute the handler with every step method short-circuited. That is plan mode.
In plan mode:
- no shell command is spawned, no HTTP request is sent, no agent is called;
- no custom
Operationis executed, so no third-party API is touched; - an approval gate is recorded and stepped over instead of suspending the run;
- a delay is recorded without sleeping;
- nothing is written: no run, no step, no dependency, no artifact.
Each step method returns a synthetic, success-shaped output instead
(exit_code: 0 for a shell step, status: 200 for an HTTP step). A handler
that branches on build.is_success() therefore follows the happy path: a plan
shows the nominal branch, not every branch the run might take.
From the CLI
ironflow run plan deploy --input '{"env":"prod"}'
workflow deploy estimated ~2m 14s
├─ build [shell] ~1m 02s
├─ parallel-1
├─ test-unit [shell] ~41s
├─ test-integration [shell] ~58s
├─ lint [shell] ~12s
└─ deploy-prod [shell] ~14s (when production run = true)
Useful flags:
| Flag | Meaning |
|---|---|
--input '<json>' | Input payload, inline. Must be a JSON object. |
--input-file <path> | Input payload read from a file. Mutually exclusive with --input. |
--max-depth <n> | How deep sub-workflows are expanded. Defaults to 3, at most 10. |
--no-estimates | Skip the duration estimate query against run history. |
--json | Emit the raw plan instead of the tree. |
From the API
POST /api/v1/workflows/{name}/plan
{
"payload": { "env": "prod" },
"max_depth": 3,
"estimate_durations": true
}
The response is the usual envelope around an execution plan:
{
"data": {
"workflow": "deploy",
"max_depth": 3,
"truncated": false,
"estimated_duration_ms": 134000,
"steps": [
{
"name": "build",
"kind": "shell",
"workflow": "deploy",
"depth": 0,
"depends_on": [],
"estimated_duration_ms": 62000
},
{
"name": "deploy-prod",
"kind": "shell",
"workflow": "deploy",
"depth": 0,
"depends_on": ["build"],
"condition": {
"state": "evaluated",
"expression": "production run",
"value": true
}
}
]
}
}
The route needs an authenticated caller, like GET /api/v1/workflows/{name}:
planning has no side effect, so it is not admin-only. It answers 400 for a
max_depth outside 1..=10 or a payload that is not a JSON object, and 404
for a workflow that is not registered.
Conditions
A condition appears on a step in one of three states:
| State | Meaning |
|---|---|
evaluated | Declared with ctx.when and resolved against the input. |
skipped | The handler called ctx.skip(name, reason) on this branch. |
unevaluable | Declared with ctx.when_dynamic: it depends on a step output, unknown before the run. |
A plain Rust if on a step output stays invisible to the planner by
construction. ctx.when_dynamic exists precisely to surface it.
Durations
When estimate_durations is on (the default), the planner samples the most
recent completed runs of the workflow and averages each step’s duration by name.
A step with no history carries no estimate, and a workflow with no history at
all reports estimated_duration_ms as absent rather than zero. The plan total
sums sequential steps and counts each parallel wave once, at its slowest member.
Sub-workflows and limits
ctx.workflow(...) is expanded in place: the child’s steps appear inline with
depth incremented and workflow set to the child’s name. max_depth bounds
that expansion, so a workflow that invokes itself stops there.
Two guards keep a plan bounded, and both set truncated with an
incomplete_reason:
- the depth limit above;
- a hard cap of 1000 planned steps, which stops a handler that loops.
A handler that returns an error mid-plan does not fail the request either: the
partial plan is returned, with incomplete_reason carrying the error.
Limits worth knowing
- Synthetic outputs are success-shaped, so only the nominal branch is planned.
- A handler that unwraps a decision answer, or deserializes
ctx.input::<T>()against a payload it does not match, aborts the plan; you get the steps recorded so far plus the reason. - Secrets and artifacts are not resolved while planning.
Transports
Transports control where agent steps execute. By default, agents run on the local machine via ClaudeCodeProvider. Ironflow ships additional transports for running agents in isolated environments.
Available transports
| Transport | Provider | Use case |
|---|---|---|
| Local | ClaudeCodeProvider | Development, simple setups |
| Docker | DockerProvider | Isolated containers on the same host |
| SSH | SshProvider | Remote machines |
| Kubernetes | K8sProvider | Ephemeral or persistent pods in a cluster |
Docker transport
Executes agent commands inside a running Docker container via docker exec:
//! Docker transport example.
//!
//! Executes `claude` inside a running Docker container via `docker exec`.
//! The container must already be running and have the `claude` binary installed.
//!
//! # Usage
//!
//! ```sh
//! DOCKER_CONTAINER=claude-worker cargo run --bin docker-transport
//! ```
use std::env;
use ironflow_core::prelude::*;
use ironflow_core::providers::claude::DockerProvider;
#[tokio::main]
async fn main() -> Result<(), OperationError> {
let container = env::var("DOCKER_CONTAINER").expect("DOCKER_CONTAINER env var required");
let provider = DockerProvider::new(&container).working_dir("/workspace");
let result = Agent::new()
.prompt("What is 2 + 2?")
.max_budget_usd(0.10)
.run(&provider)
.await?;
println!("Response: {}", result.text());
println!("Model: {}", result.model().unwrap_or("unknown"));
println!("Duration: {}ms", result.duration_ms());
Ok(())
}
SSH transport
Connects to a remote host via SSH:
//! SSH transport example.
//!
//! Connects to a remote host via SSH and runs `claude` there.
//! Requires the `claude` binary to be installed on the remote host.
//!
//! # Usage
//!
//! ```sh
//! SSH_HOST=build-server SSH_USER=deploy SSH_PASSWORD=secret cargo run --bin ssh-transport
//! ```
use std::env;
use ironflow_core::prelude::*;
use ironflow_core::providers::claude::SshProvider;
use ironflow_core::providers::claude::ssh::HostKeyPolicy;
#[tokio::main]
async fn main() -> Result<(), OperationError> {
let host = env::var("SSH_HOST").expect("SSH_HOST env var required");
let user = env::var("SSH_USER").expect("SSH_USER env var required");
let password = env::var("SSH_PASSWORD").expect("SSH_PASSWORD env var required");
// For production, use HostKeyPolicy::Fingerprint or HostKeyPolicy::KnownHostsFile
let provider = SshProvider::new(&host, &user)
.password(&password)
.host_key_policy(HostKeyPolicy::AcceptAll);
let result = Agent::new()
.prompt("What is 2 + 2?")
.max_budget_usd(0.10)
.run(&provider)
.await?;
println!("Response: {}", result.text());
println!("Model: {}", result.model().unwrap_or("unknown"));
println!("Duration: {}ms", result.duration_ms());
Ok(())
}
Kubernetes transport
Two modes are available:
- Ephemeral – creates a pod for each agent call, deletes it when done
- Persistent – reuses a long-lived pod for multiple calls
See the examples/transports/ directory for complete Kubernetes examples.
For untrusted prompts or multi-tenant clusters, K8sEphemeralProvider::sandboxed
runs each agent in a hardened pod: non-root, read-only root filesystem, secrets
read from Kubernetes Secrets, managed-settings presets and egress profiles. See
Kubernetes Sandbox.
Provider Account credentials
The local, Docker and SSH transports inject the credential of the Provider Account the worker selected: through the process environment, the Docker exec environment, or the first line of stdin over SSH. The token never appears on a command line. The Kubernetes transports do not inject accounts yet and keep using the pod environment.
Choosing a transport
- Development: use
ClaudeCodeProvider(local). No setup needed. - CI/CD pipelines: Docker or Kubernetes for isolation.
- Remote build servers: SSH for machines you already manage.
- Multi-tenant production: Kubernetes ephemeral pods for strong isolation between tenants.
Kubernetes Sandbox
K8sEphemeralProvider::sandboxed(image) runs every agent step in a hardened,
single-use pod. Everything is opt-in: K8sEphemeralProvider::new(image) keeps
its behaviour, apart from the run/step labels, the expiry annotation and the
cleanup of a previous attempt described below.
The image
The official image
registry.gitlab.com/thomastartrau/ironflow/ironflow-claude-runner:<claude-code-version>-<n>
is built from docker/claude-runner/ by the build-claude-runner-image CI job.
The current tag is in docker/claude-runner/IMAGE_TAG. There is no latest
tag: pin the full tag.
| Path | Content |
|---|---|
/home/claude | HOME of uid 10001, an emptyDir in the sandbox |
/tmp | TMPDIR, an emptyDir in the sandbox |
/etc/claude-code | managed-settings.json (baked default, replaced by a preset) |
/etc/ironflow/claude-profile/<n> | Claude profile ConfigMaps, copied into ~/.claude/<subdir> |
# Official ironflow Claude Code runner image, used by
# K8sEphemeralProvider::sandboxed.
#
# Published as ironflow-claude-runner:<claude-code-version>-<n>, the content of
# IMAGE_TAG next to this file. Never `latest`. CLAUDE_CODE_VERSION is the
# <claude-code-version> part: see README.md for a local build.
#
# Layout:
# /home/claude HOME of uid 10001. The sandboxed provider mounts
# an emptyDir here: the root filesystem is read-only.
# /tmp emptyDir mounted by the sandboxed provider.
# /etc/claude-code managed-settings.json, read by Claude Code on
# Linux. A default is baked in; a managed-settings
# ConfigMap mounted by the provider replaces it.
# /etc/ironflow/claude-profile Claude profile ConfigMap (CLAUDE.md, settings,
# agents), copied into ~/.claude before the agent
# starts because ~/.claude must stay writable.
#
# Both /etc directories are owned by root with mode 0755: the agent can read
# them, never change them.
# Pinned by digest so a rebuild of the same tag starts from the same base. To
# move to a newer node image, replace the digest and bump <n> in IMAGE_TAG.
FROM node:22-bookworm-slim@sha256:43ac6c60b8f89723f746e8a92ce91abd5017e627ce1ddfe4238355d3a30b772c
ARG CLAUDE_CODE_VERSION
RUN apt-get update \
&& apt-get install -y --no-install-recommends ca-certificates git curl \
&& rm -rf /var/lib/apt/lists/*
# Fails the build when the version is missing, or when the installed CLI
# reports another one.
RUN : "${CLAUDE_CODE_VERSION:?pass --build-arg CLAUDE_CODE_VERSION=<version>}" \
&& npm install -g "@anthropic-ai/claude-code@${CLAUDE_CODE_VERSION}" \
&& npm cache clean --force \
&& test "$(claude --version | cut -d' ' -f1)" = "${CLAUDE_CODE_VERSION}"
RUN groupadd -g 10001 claude \
&& useradd -u 10001 -g 10001 -m -d /home/claude claude
# Marketplace neutralisation: no plugin marketplace, no hooks, no auto-update,
# no non-essential traffic. A mounted managed-settings ConfigMap replaces the
# baked file.
RUN mkdir -p /etc/claude-code /etc/ironflow/claude-profile \
&& printf '%s\n' '{"strictKnownMarketplaces": [], "disableAllHooks": true}' \
> /etc/claude-code/managed-settings.json \
&& chown -R root:root /etc/claude-code /etc/ironflow \
&& chmod 0755 /etc/claude-code /etc/ironflow /etc/ironflow/claude-profile \
&& chmod 0644 /etc/claude-code/managed-settings.json
ENV CLAUDE_CODE_DISABLE_NONESSENTIAL_TRAFFIC=1 \
DISABLE_AUTOUPDATER=1 \
HOME=/home/claude
USER 10001:10001
WORKDIR /home/claude
sandboxed() defaults
| Default | Relaxation |
|---|---|
Runs as uid/gid 10001, runAsNonRoot, seccomp RuntimeDefault | .run_as_user(uid) (never 0) |
| Read-only root filesystem, all capabilities dropped, no privilege escalation | .allow_writable_root() |
HOME on a 1Gi emptyDir | .home_size_limit("4Gi") |
/tmp on a 512Mi emptyDir | .tmp_size_limit("2Gi") |
activeDeadlineSeconds = timeout + 60s | .deadline_margin(d), or .active_deadline_seconds(d) to set it outright |
| No service account token mounted | .service_account(name) on the provider or the step |
| Secrets refused as plain text | none: use a Secret |
A relaxation called on a provider built with new() panics.
Secrets
A sandboxed provider refuses oauth_credentials(json) and
env("ANTHROPIC_API_KEY", ..) / env("CLAUDE_CODE_OAUTH_TOKEN", ..): the
value would sit in clear text in the pod spec. Read it from a Kubernetes Secret
instead; the pod only carries a secretKeyRef:
let provider = K8sEphemeralProvider::sandboxed(&image)
// Long-lived token from `claude setup-token`.
.oauth_token_from_secret("claude-oauth", "token")
// Or the full credentials JSON, written to ~/.claude/.credentials.json.
.oauth_credentials_from_secret("claude-credentials", "credentials.json")
// Any other variable.
.env_from_secret("GITLAB_TOKEN", "gitlab-bot", "token");
Auth proxy: no Claude credential in the pod
A secretKeyRef keeps the credential out of the pod spec, not out of the pod:
anything running in the agent container can still read
CLAUDE_CODE_OAUTH_TOKEN. With auth_proxy, the pod never receives it:
let provider = K8sEphemeralProvider::sandboxed(&image)
.namespace("ironflow-agents")
.auth_proxy("http://ironflow-auth-proxy.ironflow-system");
The ironflow-auth-proxy service (crate ironflow-auth-proxy, manifests in
examples/k8s/sandbox/auth-proxy.yaml) holds the credential instead. Its
official image
registry.gitlab.com/thomastartrau/ironflow/ironflow-auth-proxy:<version> is
built from docker/auth-proxy/Dockerfile by the build-auth-proxy-image CI
job. <version> is the version of the ironflow-auth-proxy crate: the job
publishes it once that version is released, and never rebuilds a published
tag. There is no latest tag: pin the version.
- Token lifecycle. At pod launch the worker calls the proxy admin API and
gets an opaque token (
ifap_...) bound to the run id, the step and the pod expiry (ironflow.io/expires-at). The pod receivesANTHROPIC_BASE_URL(the proxy URL),ANTHROPIC_AUTH_TOKEN(the opaque token) andCLAUDE_CODE_DISABLE_NONESSENTIAL_TRAFFIC=1. The worker revokes the token when the step ends, whatever the outcome, and every token of a run when the run is released; the expiry is the backstop. An unknown, revoked or expired token gets a 401. If the proxy cannot issue a token, the step fails before any pod is created. - Restrictions. The proxy only relays GET and POST requests under
/v1/toapi.anthropic.com. A path outside/v1/(including/admin, which needs the admin key) or a request for another host (absolute-form URI,CONNECT) gets a 403, any other method a 405. - Which credential. The step’s Provider Account when one is attached (the
provider reports the
claude_subscriptionaccount kind onceauth_proxyis set), elseCLAUDE_CODE_OAUTH_TOKEN, thenANTHROPIC_API_KEY, from the worker environment. The worker also needsIRONFLOW_AUTH_PROXY_ADMIN_KEY(or.auth_proxy_admin_key(..)). Rate-limit windows reported through the proxy are recorded on the Provider Account as with the Docker provider. - No credential on the pod side. With
auth_proxyset, a Claude credential configured for the pod (oauth_token_from_secret,oauth_credentials,oauth_credentials_from_secret, a stepenv_from_secret("CLAUDE_CODE_OAUTH_TOKEN", ..), any plain value starting withsk-ant) fails the step with an error naming the variable. - Registry. In memory by default (one replica). Set
IRONFLOW_AUTH_PROXY_DATABASE_URLand an encryption key (IRONFLOW_SECRET_KEYS) for a shared PostgreSQL registry: several replicas, tokens survive restarts, the proxy needs egress to PostgreSQL. Under Cilium, uncomment the 5432 rule in section (b) ofcilium-egress-auth-proxy.yaml(below). With standard NetworkPolicies only, nothing restricts the proxy’s egress unless the operator adds a policy on the proxy pods. - Logs. The proxy and the worker log the first 12 characters of the token id (a SHA-256 of the token), never the token or the credential.
The network side changes too: the agent pods only reach the proxy, and the
proxy only reaches api.anthropic.com (and PostgreSQL with the shared
registry):
# Cilium egress for the auth proxy variant.
#
# Use INSTEAD of cilium-egress-anthropic.yaml: the agent pods no longer reach
# api.anthropic.com, only the auth proxy, which alone holds the credential
# and alone reaches the API. The per-profile policies (gitlab below) are
# unchanged and stay additive.
#
# (a) Agent pods: DNS, and the auth proxy pods on 8080.
apiVersion: cilium.io/v2
kind: CiliumNetworkPolicy
metadata:
name: claude-runner-egress-auth-proxy
namespace: ironflow-agents
spec:
endpointSelector:
matchLabels:
app.kubernetes.io/component: claude-runner
egress:
- toEndpoints:
- matchLabels:
k8s:io.kubernetes.pod.namespace: kube-system
k8s:k8s-app: kube-dns
toPorts:
- ports:
- port: "53"
protocol: ANY
rules:
dns:
- matchPattern: "*"
- toEndpoints:
- matchLabels:
k8s:io.kubernetes.pod.namespace: ironflow-system
app.kubernetes.io/name: ironflow-auth-proxy
toPorts:
- ports:
- port: "8080"
protocol: TCP
---
# Per-profile opening, as in cilium-egress-anthropic.yaml: pods labelled
# ironflow.io/egress-profile=gitlab may also reach gitlab.com.
apiVersion: cilium.io/v2
kind: CiliumNetworkPolicy
metadata:
name: claude-runner-egress-gitlab
namespace: ironflow-agents
spec:
endpointSelector:
matchLabels:
app.kubernetes.io/component: claude-runner
ironflow.io/egress-profile: gitlab
egress:
- toEndpoints:
- matchLabels:
k8s:io.kubernetes.pod.namespace: kube-system
k8s:k8s-app: kube-dns
toPorts:
- ports:
- port: "53"
protocol: ANY
rules:
dns:
- matchPattern: "*"
- toFQDNs:
- matchName: gitlab.com
toPorts:
- ports:
- port: "443"
protocol: TCP
---
# (b) The auth proxy pods: egress to DNS and api.anthropic.com:443 only, plus
# PostgreSQL on 5432 with the shared registry (commented-out rule below);
# ingress on 8080 only from the agent pods and from the worker. Adjust the
# worker selector to the namespace and labels the worker runs under.
apiVersion: cilium.io/v2
kind: CiliumNetworkPolicy
metadata:
name: ironflow-auth-proxy
namespace: ironflow-system
spec:
endpointSelector:
matchLabels:
app.kubernetes.io/name: ironflow-auth-proxy
ingress:
- fromEndpoints:
- matchLabels:
k8s:io.kubernetes.pod.namespace: ironflow-agents
app.kubernetes.io/component: claude-runner
- matchLabels:
k8s:io.kubernetes.pod.namespace: ironflow
app.kubernetes.io/name: ironflow-worker
toPorts:
- ports:
- port: "8080"
protocol: TCP
egress:
- toEndpoints:
- matchLabels:
k8s:io.kubernetes.pod.namespace: kube-system
k8s:k8s-app: kube-dns
toPorts:
- ports:
- port: "53"
protocol: ANY
rules:
dns:
- matchPattern: "*"
- toFQDNs:
- matchName: api.anthropic.com
toPorts:
- ports:
- port: "443"
protocol: TCP
# Shared PostgreSQL registry (IRONFLOW_AUTH_PROXY_DATABASE_URL set in
# auth-proxy.yaml): uncomment ONE of the two rules below, else the proxy
# times out at startup. Not needed with the in-memory registry. If the
# database has its own ingress policy, open 5432 there to the proxy pods.
#
# In-cluster database: adjust its namespace and pod labels.
# - toEndpoints:
# - matchLabels:
# k8s:io.kubernetes.pod.namespace: ironflow-db
# app.kubernetes.io/name: postgresql
# toPorts:
# - ports:
# - port: "5432"
# protocol: TCP
#
# Managed database outside the cluster: adjust to its subnet.
# - toCIDR:
# - 10.0.0.0/24
# toPorts:
# - ports:
# - port: "5432"
# protocol: TCP
Per-step settings
A step adds to or overrides the provider’s settings through its
AgentConfig. Other providers ignore these fields.
let config = AgentConfig::new("Open the merge request")
.env_from_secret("GITLAB_TOKEN", "gitlab-bot", "token") // wins over the provider's
.service_account("gitlab-reader") // wins over the provider's
.read_only_pvc("repos", "/data/repos") // after the provider's volumes
.read_only_config_map("guidelines", "/data/guidelines")
.managed_settings("readonly") // preset registered on the provider
.runtime_class("gvisor") // wins over the provider's
.egress_profile("gitlab"); // ironflow.io/egress-profile label
Read-only mounts must be absolute, unique, and cannot shadow /,
/home/claude, /tmp, /etc/claude-code, or /etc/ironflow/claude-profile
and anything below it.
Claude profile
A Claude profile (CLAUDE.md, settings.json, rules/, agents/,
commands/) reaches ~/.claude through ConfigMaps. A ConfigMap key cannot
contain /, so each directory of the profile is its own ConfigMap, mapped to
its sub-directory:
let provider = K8sEphemeralProvider::sandboxed(&image)
.claude_profile_configmap("claude-profile") // ~/.claude
.claude_profile_configmap_at("claude-profile-rules", "rules") // ~/.claude/rules
.claude_profile_configmap_at("claude-profile-agents", "agents");
Each ConfigMap is mounted read-only in its own directory, then only its keys
are copied into ~/.claude/<subdir> before the agent starts: never the
..data entries of the volume. Profiles are copied before the credentials,
which they cannot overwrite. A subdir must be relative, made of
[A-Za-z0-9._-] segments without . or .., and mapped once: the provider
panics at build time otherwise.
With kustomize, one configMapGenerator per directory. Disable the name hash:
the provider refers to the ConfigMaps by name.
# kustomization.yaml, next to claude-home/
namespace: ironflow-agents
generatorOptions:
disableNameSuffixHash: true
configMapGenerator:
- name: claude-profile
files:
- claude-home/CLAUDE.md
- claude-home/settings.json
- name: claude-profile-rules
files:
- claude-home/rules/rust.md
- claude-home/rules/security.md
kustomize lists every file. To pick up a whole directory instead:
kubectl create configmap claude-profile-rules --from-file=claude-home/rules/ --dry-run=client -o yaml.
Managed-settings presets map a name to a ConfigMap holding
managed-settings.json. An unknown preset fails the step, never falls back:
let provider = K8sEphemeralProvider::sandboxed(&image)
.managed_settings_preset("locked", "claude-managed-locked")
.managed_settings_preset("readonly", "claude-managed-readonly")
.default_managed_settings("locked");
Labels and retry cleanup
The engine stamps ironflow.io/run-id, ironflow.io/root-run-id and
ironflow.io/step on every agent step (outside the engine, call
AgentConfig::run_scope(run_id, step)). Step names are sanitized into valid
label values, with a hash suffix when altered. The root run is the run itself,
or the top-level run inside a sub-workflow.
Before every execution of a run, first one included, the engine calls
release_run: the provider deletes every pod, JobRun Job and prompt
ConfigMap labelled with the run id or with the run as root, and waits until
the pods are gone (.previous_attempt_timeout(d), 60s by default). A retry
that starts by resetting shared state (a worktree) never runs next to an
agent of the dead attempt still writing to it. If the pods are still there, or
the Kubernetes API fails, the execution fails with a replayable error and the
next attempt tries again.
Tag a pod you create yourself with the same labels so that it is released too:
let run = PodRun::new(&kube, "check", &image, "cargo test")
.label(LABEL_RUN_ID, &ctx.run_id().to_string())
.label(LABEL_ROOT_RUN_ID, &ctx.root_run_id().to_string());
Before creating a pod, the provider also deletes the pods and prompt
ConfigMaps of a previous attempt of the same step of the same run. If they are
still terminating, the step fails: two agents never run side by side. Two
branches of a ctx.parallel() group with the same name would delete each
other’s pod, so the engine fails such a group before creating any step.
Deleting JobRun Jobs needs list and delete on jobs; without them, Jobs
are skipped with a warning and the pods are still released.
The reaper
Every object ironflow creates carries app.kubernetes.io/managed-by=ironflow,
an app.kubernetes.io/component (claude-runner, prompt-data, pod-run,
job-run) and the ironflow.io/expires-at annotation (unix seconds: creation
- timeout + margin). That covers the agent pods and their prompt ConfigMaps,
and the pods of
PodRunand Jobs ofJobRunfromironflow-ops-k8s(.expiry_margin(d), 60s by default). A caller cannot setmanaged-byorcomponent:pod_label,PodRun::labelandJobRun::labelpanic, an agent step carrying one fails.
The reaper selects on managed-by=ironflow and deletes pods and Jobs past
their expiry or killed with DeadlineExceeded, and expired prompt ConfigMaps.
A Job goes with its pods; a pod a Job controls is left to its Job. Objects
without a parseable annotation are never touched.
let report = provider.reap_orphans().await?; // one pass
let handle = provider.spawn_orphan_reaper(Duration::from_secs(300)); // background
// Without a K8sEphemeralProvider, e.g. a worker that only runs PodRun:
let report = reap_orphans(&K8sClusterConfig::Default, "ironflow-agents").await?;
Reaping Jobs needs list and delete on jobs (see
examples/k8s/sandbox/namespace-rbac.yaml). Without them, the Job pass is
skipped with a warning and the rest of the pass runs.
gVisor (RuntimeClass)
A runc container shares the node kernel: a kernel exploit in untrusted code
reaches the node. gVisor runs the pod against a user-space kernel instead. On
Talos, install the gvisor system extension on the nodes, then declare the
RuntimeClass once per cluster:
apiVersion: node.k8s.io/v1
kind: RuntimeClass
metadata:
name: gvisor
handler: runsc
Set it on the provider to cover every agent pod, and override it for one step (the step wins):
let provider = K8sEphemeralProvider::sandboxed(&image)
.runtime_class("gvisor");
let config = AgentConfig::new("Review this untrusted patch")
.runtime_class("kata"); // wins over the provider's
A check pod created with PodRun takes the same setting:
let run = PodRun::new(&kube, "check", "rust:1.94", "cargo test")
.runtime_class("gvisor");
Without a value, spec.runtimeClassName stays absent and the cluster default
runtime (runc) applies. A blank name is refused.
Warning: builds are noticeably slower under gVisor (compilation, many small file syscalls,
cargoandnpminstalls). Keep gVisor for the steps that handle untrusted code and leave trusted build steps onrunc.
Network policies
The worker does not create network policies: its Role in
examples/k8s/sandbox/namespace-rbac.yaml has no access to them. A cluster
administrator applies a default deny (networkpolicy-deny-all.yaml) and FQDN
openings selected on the ironflow.io/egress-profile label:
# FQDN egress for agent pods with Cilium.
#
# Every claude-runner pod may resolve names through kube-dns and reach
# api.anthropic.com on 443. The DNS rule with `matchPattern: "*"` is required:
# Cilium learns the IPs behind toFQDNs names by inspecting DNS answers.
apiVersion: cilium.io/v2
kind: CiliumNetworkPolicy
metadata:
name: claude-runner-egress-anthropic
namespace: ironflow-agents
spec:
endpointSelector:
matchLabels:
app.kubernetes.io/component: claude-runner
egress:
- toEndpoints:
- matchLabels:
k8s:io.kubernetes.pod.namespace: kube-system
k8s:k8s-app: kube-dns
toPorts:
- ports:
- port: "53"
protocol: ANY
rules:
dns:
- matchPattern: "*"
- toFQDNs:
- matchName: api.anthropic.com
toPorts:
- ports:
- port: "443"
protocol: TCP
---
# Per-profile opening: pods labelled ironflow.io/egress-profile=gitlab (set
# with `.egress_profile("gitlab")` on the provider or the step) may also reach
# gitlab.com. Cilium policies are additive: this adds to the rule above.
apiVersion: cilium.io/v2
kind: CiliumNetworkPolicy
metadata:
name: claude-runner-egress-gitlab
namespace: ironflow-agents
spec:
endpointSelector:
matchLabels:
app.kubernetes.io/component: claude-runner
ironflow.io/egress-profile: gitlab
egress:
- toEndpoints:
- matchLabels:
k8s:io.kubernetes.pod.namespace: kube-system
k8s:k8s-app: kube-dns
toPorts:
- ports:
- port: "53"
protocol: ANY
rules:
dns:
- matchPattern: "*"
- toFQDNs:
- matchName: gitlab.com
toPorts:
- ports:
- port: "443"
protocol: TCP
Full example
//! Sandboxed K8s ephemeral transport example.
//!
//! Runs Claude Code in a hardened pod: non-root, read-only root filesystem,
//! all capabilities dropped, OAuth token read from a Kubernetes Secret, a
//! locked-down managed-settings preset, and the `anthropic-only` egress
//! profile label. A background task reaps orphaned pods every five minutes.
//!
//! Apply the manifests of `examples/k8s/sandbox/` first. Requires a reachable
//! Kubernetes cluster (via kubeconfig or in-cluster).
//!
//! # Usage
//!
//! ```sh
//! K8S_IMAGE=registry.gitlab.com/thomastartrau/ironflow/ironflow-claude-runner:2.1.284-1 \
//! cargo run --bin k8s-sandboxed
//! ```
use std::env;
use std::time::Duration;
use ironflow_core::prelude::*;
use ironflow_core::providers::claude::K8sEphemeralProvider;
#[tokio::main]
async fn main() -> Result<(), OperationError> {
let image = env::var("K8S_IMAGE").expect("K8S_IMAGE env var required");
let provider = K8sEphemeralProvider::sandboxed(&image)
.namespace("ironflow-agents")
.oauth_token_from_secret("claude-oauth", "token")
.managed_settings_preset("locked", "claude-managed-locked")
.default_managed_settings("locked")
.egress_profile("anthropic-only")
.timeout(Duration::from_secs(600));
let _reaper = provider.spawn_orphan_reaper(Duration::from_secs(300));
// Outside the engine, `run_scope` sets the run/step labels the engine
// stamps on every agent step.
let config = AgentConfig::new("List the top-level directories of /data/repos.")
.max_budget_usd(0.10)
.read_only_pvc("repos", "/data/repos")
.run_scope("demo-run", "investigate");
let result = Agent::from_config(config).run(&provider).await?;
println!("Response: {}", result.text());
println!("Model: {}", result.model().unwrap_or("unknown"));
println!("Duration: {}ms", result.duration_ms());
Ok(())
}
Checks
Run these against a live agent pod before trusting the setup:
# Non-root, read-only root filesystem, capabilities dropped.
kubectl -n ironflow-agents get pod -l app.kubernetes.io/component=claude-runner \
-o jsonpath='{.items[0].spec.containers[0].securityContext}'
# No secret value in the spec, only secretKeyRef.
kubectl -n ironflow-agents get pod <pod> -o yaml | grep -A3 CLAUDE_CODE_OAUTH_TOKEN
# No service account token mounted.
kubectl -n ironflow-agents exec <pod> -- ls /var/run/secrets/kubernetes.io/serviceaccount
# Root filesystem is read-only.
kubectl -n ironflow-agents exec <pod> -- touch /usr/local/probe
# Egress is limited to the profile.
kubectl -n ironflow-agents exec <pod> -- curl -sS -m 5 https://example.com
# The worker cannot touch network policies.
kubectl auth can-i create networkpolicies -n ironflow-agents \
--as=system:serviceaccount:ironflow:ironflow-worker
# Run, root run and step labels, and the expiry annotation.
kubectl -n ironflow-agents get pods \
-L ironflow.io/run-id,ironflow.io/root-run-id,ironflow.io/step \
-o custom-columns='NAME:.metadata.name,EXPIRES:.metadata.annotations.ironflow\.io/expires-at'
The exec checks against the service account, the root filesystem and egress
must fail; kubectl auth can-i must answer no.
With the auth proxy, also check that the pod holds no Claude secret and that the proxy refuses an unknown token:
kubectl -n ironflow-agents exec <pod> -- env | grep -c 'sk-ant' # 0
kubectl -n ironflow-agents exec <pod> -- sh -c \
'curl -s -o /dev/null -w "%{http_code}" -H "Authorization: Bearer invalide" "$ANTHROPIC_BASE_URL/v1/messages"' # 401
Testing Workflows
ironflow_engine::testing::TestEngine runs a WorkflowHandler against an
in-memory store with mocked steps. The run, the steps, the FSM transitions and
the persistence are the production ones – only the outside world is swapped
out.
What it replaces:
| Production | Under TestEngine |
|---|---|
| API server | nothing to start; the run executes inline |
| Background worker | nothing to start; run() returns once the run is finished |
| Postgres | InMemoryStore |
sh -c <command> | a closure |
| An HTTP request | a closure |
| The Claude CLI | a closure, or a recorded fixture |
| A human clicking Approve | an ApprovalOutcome |
A first test
use ironflow_engine::prelude::*;
use ironflow_engine::testing::{MockShellOutput, TestEngine};
use ironflow_store::models::{RunStatus, StepStatus};
use serde_json::json;
use crate::handlers::Deploy;
#[tokio::test]
async fn deploy_runs_build_then_ship() {
let result = TestEngine::new()
.with_handler(Deploy)
.with_mock_shell(|cfg| match cfg.command.as_str() {
"cargo build" => Ok(MockShellOutput::ok("compiled")),
_ => Ok(MockShellOutput::ok("shipped")),
})
.run(json!({"environment": "staging"}))
.await
.expect("the harness ran the handler");
assert_eq!(result.status(), RunStatus::Completed);
assert_eq!(result.step_names(), vec!["build", "ship"]);
assert_eq!(result.step("build").step_output().stdout(), "compiled");
assert_eq!(result.step("ship").status(), StepStatus::Completed);
}
A handler that fails is not an Err: the returned TestResult carries
RunStatus::Failed and the message in error(). Only wiring failures – no
handler registered, two handlers sharing a name, a store rejection – come back
as Err.
Building the harness
| Method | What it does |
|---|---|
with_handler(handler) | Registers a handler. The first one is what run() executes. |
with_mock_shell(f) | Answers every shell step from f(&ShellConfig). |
with_mock_http(f) | Answers every HTTP step from f(&HttpConfig). |
with_mock_agent(f) | Answers every agent step from f(&AgentConfig). |
with_recorded_agent(dir) | Replays agent fixtures from dir. |
with_agent_provider(p) | Uses an arbitrary AgentProvider. |
with_decision_provider(p) | Wires a DecisionProvider for ctx.decision(...). |
with_mock_approval(outcome) | Resolves every approval gate with outcome. |
with_secret(key, value) | Seeds a secret for the workflow under test (secret-store feature). |
store() | The InMemoryStore, for assertions the accessors do not cover. |
Every with_* method panics if called after the first run: the engine is built
once, so a later change would be silently ignored.
Then run:
| Method | What it does |
|---|---|
run(payload) | Runs the first registered handler. |
run_workflow(name, payload) | Runs a specific registered handler. |
resume(run_id) | Continues a run suspended on an approval gate. |
Asserting on the result
TestResult reads the run and its steps back from the store, so an assertion
sees exactly what the API and the dashboard would serve.
| Accessor | Returns |
|---|---|
status() | The RunStatus the run finished in. |
is_completed() | Whether that status is Completed. |
error() | Why the run stopped, if it did not complete. |
steps() | Every persisted step, ordered by position. |
step_names() | Those steps’ names, in the same order. |
step(name) | The first step with that name; panics when there is none. |
try_step(name) | The same, as an Option. |
output() | The last step’s output. |
duration(), cost_usd() | The run totals. |
run_id(), run() | The run identity and the raw record. |
step_results() | Per-step metrics, empty when the run failed. |
Each TestStep exposes name(), kind(), status(), step_output(),
output(), input(), error(), duration(), cost_usd(), is_completed(),
is_error_handler() and raw(). step_output() reads the persisted output
through the typed StepOutput accessors (stdout(), status(), body(),
text(), json::<T>()); output() is the raw JSON.
Steps of a parallel wave share a position and a handler may reuse a name:
disambiguate those with steps() rather than step(name).
Shell and HTTP parity
The mocks reproduce the asymmetry of the real executors, so allow_failure,
step retries and run failure behave exactly as in production:
- A
MockShellOutputwith a non-zeroexit_codeis an error, like a real non-zero exit. UseMockShellOutput::failed(1, "boom"). When the step setsexit_code_as_output(), the mock completes with that exit code as its output. - A
MockHttpResponsewith a non-2xxstatusis a normal output, like a real 500 response. ReturnErr(OperationError::Http { status: None, .. })from the closure to simulate a transport failure instead.
use ironflow_core::error::OperationError;
use ironflow_engine::testing::{MockHttpResponse, TestEngine};
use serde_json::json;
let harness = TestEngine::new()
.with_handler(Fetch)
// A 404 the handler is expected to deal with.
.with_mock_http(|cfg| {
if cfg.url.ends_with("/missing") {
Ok(MockHttpResponse::json(404, &json!({"error": "not found"})))
} else {
Ok(MockHttpResponse::ok(&json!({"id": 7})))
}
});
Approval gates
Two ways to test a gated handler:
use ironflow_engine::testing::{ApprovalOutcome, TestEngine};
// 1. Resolve the gate inline and assert on the whole run.
let approved = TestEngine::new()
.with_handler(GatedDeploy)
.with_mock_shell(|_cfg| Ok(MockShellOutput::ok("ok")))
.with_mock_approval(ApprovalOutcome::Approved)
.run(json!({}))
.await?;
assert_eq!(approved.status(), RunStatus::Completed);
// A rejection fails the run with EngineError::ApprovalRejected.
let rejected = TestEngine::new()
.with_handler(GatedDeploy)
.with_mock_shell(|_cfg| Ok(MockShellOutput::ok("ok")))
.with_mock_approval(ApprovalOutcome::reject("budget freeze"))
.run(json!({}))
.await?;
assert_eq!(rejected.step("gate").status(), StepStatus::Rejected);
// 2. Without a mock, the gate suspends the run, the way production does.
let mut harness = TestEngine::new()
.with_handler(GatedDeploy)
.with_mock_shell(|_cfg| Ok(MockShellOutput::ok("ok")));
let suspended = harness.run(json!({})).await?;
assert_eq!(suspended.status(), RunStatus::AwaitingApproval);
let resumed = harness.resume(suspended.run_id()).await?;
assert_eq!(resumed.status(), RunStatus::Completed);
Agent fixtures
with_recorded_agent(dir) replays fixtures written by RecordReplayProvider.
The argument is the directory: each fixture is keyed by a hash of the
AgentConfig and stored as <hash>.json inside it. A missing fixture fails the
step instead of falling back to the real Claude CLI, so a stale suite never
silently starts spending tokens.
let result = TestEngine::new()
.with_handler(Review)
.with_recorded_agent("tests/fixtures")
.run(json!({}))
.await?;
To record, pass a recording provider through the escape hatch:
use std::sync::Arc;
use ironflow_core::provider::AgentProvider;
use ironflow_core::providers::claude::ClaudeCodeProvider;
use ironflow_core::providers::record_replay::RecordReplayProvider;
let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::record(
ClaudeCodeProvider::new(),
"tests/fixtures",
));
let result = TestEngine::new()
.with_handler(Review)
.with_agent_provider(provider)
.run(json!({}))
.await?;
With no agent backend configured at all, an agent step fails with a message naming the three constructors – a forgotten mock is a loud failure, not a network call.
Parallel waves, error handlers and sub-workflows
The mocks apply to the steps inside them: a step of a ctx.parallel(...) wave,
a step fired by ctx.on_error(...), and every step of a child run started with
ctx.workflow(...) all go through the same interceptor. Register both handlers
and drive the parent:
let mut harness = TestEngine::new()
.with_handler(Parent)
.with_handler(Child)
.with_mock_shell(|_cfg| Ok(MockShellOutput::ok("ok")));
let store = harness.store();
let result = harness.run_workflow("parent", json!({})).await?;
// A workflow step stores a `SubWorkflowOutput`: read it back typed.
let child: SubWorkflowOutput = result.step("child").step_output().json()?;
let child_steps = store.list_steps(child.run_id()).await?;
Limitations
- Custom operations (
ctx.operation(...)) are not intercepted. Mock one by passing a test-doubleOperationto the handler. ctx.delay(...)is not intercepted: a non-zero delay still suspends the run withRunStatus::Sleeping. It resumes onceRunWaker::tickruns after itsscheduled_at.ctx.wait_for_signal(...)is resolved withwith_mock_signal(|step, name, key| ..), returningSignalOutcome::Received(json!(..))orSignalOutcome::TimedOut. Without it, the run ends inRunStatus::Sleepinguntil a signal is delivered.ctx.decision(...)needs a realDecisionProvider, wired withwith_decision_provider.
Use the real Engine when the test must exercise real commands, real requests
or a real agent; use TestEngine when it must exercise the handler’s logic.
Architecture Overview
Ironflow follows a client-server architecture with background workers for execution.
Components
graph TD
Dashboard[Web Dashboard] --> API[API Server]
CLI[CLI] --> SDK[Rust SDK]
MCP[MCP Server] --> SDK
SDK --> API
API --> Store[(Database)]
API --> Artifacts[(Blob Store)]
Worker1[Worker 1] --> API
Worker2[Worker 2] --> API
Worker1 --> Provider[Agent Provider]
Worker2 --> Provider
Crate map
| Crate | Role |
|---|---|
ironflow-core | Shell execution, agent providers, cost tracking |
ironflow-engine | Workflow handler trait, context, step orchestration |
ironflow-api | REST API (axum), routes, SSE events, dashboard serving |
ironflow-worker | Background worker that polls and executes runs |
ironflow-store | Storage trait + PostgreSQL and in-memory backends |
ironflow-auth | JWT authentication, password hashing, API keys |
ironflow-runtime | Daemon features: webhooks, trigger sources |
ironflow-artifacts | Blob storage for step-produced files |
ironflow-templates | Fetch and install workflow templates from Git |
ironflow-sdk | Type-safe Rust client (types generated from OpenAPI) |
ironflow-cli | Command-line interface (clap v4) |
ironflow-mcp | Model Context Protocol server |
ironflow-types | Shared API envelope types |
Request flow
- A client (dashboard, CLI, SDK, or webhook) sends a request to the API
- The API validates authentication, creates a Run in the Store, and returns it
- A Worker polls the API, acquires a lease on the Run, and executes the handler
- The handler calls steps (
ctx.shell(),ctx.agent(), etc.), each persisted as they complete - Events are published via SSE for real-time updates
- On completion or failure, the Worker reports the result back to the API
Data flow
Runs and steps are stored in PostgreSQL (or in-memory for development). Artifacts (files produced by steps) are stored in a separate blob store (local filesystem by default). The two are linked by artifact metadata on each step.