Keyboard shortcuts

Press ← or → to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

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.

Installation

Prerequisites

  • Rust 1.94+ (see rust-version in 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:

  1. A library crate with your workflow handlers
  2. A server binary that exposes the API and serves the dashboard
  3. 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

VariableDefaultDescription
IRONFLOW_ENVdevelopmentproduction or development
DATABASE_URL–PostgreSQL URL (required in production)
JWT_SECRETdev secretJWT signing key (in production: mandatory, >= 32 bytes, must not start with ironflow-dev-)
WORKER_TOKENdev tokenShared secret for worker auth (in production: mandatory, >= 32 bytes, must not start with ironflow-dev-)
PORT3000HTTP listen port
ALLOWED_ORIGINSsame-originComma-separated CORS origins
ARTIFACTS_DIR–Filesystem root for step artifacts
PURGE_MAX_AGE_DAYS90Terminal runs older than this are purged
PURGE_MAX_RUNS_PER_WORKFLOW1000Terminal runs kept per workflow
PURGE_DRY_RUNfalseLog what would be purged without deleting (also disables usage and signal purging)
PURGE_INTERVAL_SECS86400Seconds between purger ticks (min 60)
PROVIDER_ACCOUNT_USAGE_RETENTION_DAYS30Days of Provider Account usage history kept (min 1)
SIGNAL_RETENTION_DAYS7Days 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

VariableDefaultDescription
API_URLhttp://localhost:3000Address of the API server
WORKER_TOKENdev tokenShared secret matching the server
CONCURRENCY2Number of parallel runs
POLL_INTERVAL_SECS2Seconds 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 a WorkflowContext to 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 through ctx.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

KindMethodDescription
Shellctx.shell()Execute a shell command
Httpctx.http()Make an HTTP request
Agentctx.agent()Call an AI agent (Claude, OpenAI, etc.)
Approvalctx.approval()Pause for human approval
HumanInputctx.human_input()Pause until a human submits a typed answer (Human Input)
Signalctx.wait_for_signal()Pause until an external signal arrives or a timeout elapses (Signals)
Decisionctx.decision()Make a typed machine decision (System One / Jev)
Workflowctx.workflow()Start a sub-workflow
Customctx.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?;
FieldBuilderMeaning
messageApprovalConfig::newPrompt shown to reviewers
assigneeassigned_toAssignee::user / Assignee::group expected to answer
deadline_secswith_deadline / with_deadline_secsSLA window, in seconds
on_timeouton_timeoutEscalationPolicy applied when the deadline fires (defaults to AutoReject)
timeout_secondswith_timeout_secondsLegacy spelling of a deadline with an implicit AutoReject
approversrequiringApprovers 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 database
  • execute() runs the operation and returns JSON output
  • input() 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 to Running and calls Engine::resume_run itself. Use it for single-process deployments where the API also registers the handlers. TestEngine always resumes this way.
  • ExecutionMode::Workers: the API moves the run back to Pending. A worker claims it through pick_next_pending and finishes it with Engine::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:

  1. Waits for a free execution slot (concurrency)
  2. Polls the API for a pending run and acquires a lease on it
  3. Executes the workflow handler via the Engine
  4. Refreshes the lease periodically during execution
  5. 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

TransportHow 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
Kubernetesnot 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:

StrategyPicks
least_utilized (default)the lowest peak utilization, plus 0.15 per running step
prioritythe lowest priority value
round_robinthe 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

  1. The handler calls ctx.approval() with a prompt message
  2. The run transitions to AwaitingApproval
  3. The worker releases the run and moves on to other work
  4. A human calls POST /api/v1/runs/:id/approve or POST /api/v1/runs/:id/reject
  5. On approval, the run resumes. Under ExecutionMode::Workers it is requeued to Pending: a worker picks it up, replays completed steps from cache, skips the approved gate, and continues execution. Under ExecutionMode::Local (the default) the API process resumes it the same way itself
  6. 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.

PolicyWhat it does when the deadline fires
AutoApproveCompletes 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).
AutoRejectFails 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_remaining and approval_assignee on every step of GET /api/v1/runs/:id (clamped at 0, null without a deadline);
  • the dashboard: a countdown badge on the gate in the run’s step list;
  • the CLI: the SLA column of ironflow 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:

EndpointWhat it does
POST /api/v1/approval-delegationsGrant a delegation. The delegator is always the caller.
GET /api/v1/approval-delegationsList 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?;
BuilderMeaning
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 /approve returns 200 with the run still awaiting_approval, the gate keeps its SLA timer, and the CLI prints Approval 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_requested is published when the gate opens and carries the recorded requirement.
  • approval_granted is published for every vote, with the step_id, approvals_received, approvals_required and the requirement. The gate resolves when approvals_received >= approvals_required.
  • approval_rejected carries the step_id and the requirement.
{
  "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

  1. The handler calls ctx.human_input::<T>() with a message. T derives Deserialize and JsonSchema.
  2. The engine records a step of kind human_input whose input holds the message and the JSON schema of T (under the key schema).
  3. The step and the run move to AwaitingApproval, the same status as an approval gate, and an input_required event is published on the run’s event stream.
  4. 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 under ExecutionMode::Local (the default), or by requeuing it to Pending for a worker under ExecutionMode::Workers (see execution mode).
  5. The handler is replayed: completed steps come from cache and human_input returns the answer, deserialized into T.

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
FieldBuilderMeaning
messageHumanInputConfig::newPrompt shown to the person answering
assigneeassigned_toAssignee::user / Assignee::group expected to answer
approversrequiringGroups allowed to answer. The first valid answer wins, whatever the count
deadline_secswith_deadline / with_deadline_secsSLA window, in seconds
on_timeouton_timeoutEscalationPolicy 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"]}'
  • 200 returns the run, now running.

  • 422 with code INVALID_INPUT when the body does not match the schema; each violation is listed in error.details.errors:

    {
      "error": {
        "code": "INVALID_INPUT",
        "message": "input does not match the expected schema",
        "details": { "errors": ["3 is not of type \"array\""] }
      }
    }
    
  • 409 when the input was already answered or rejected.

  • 400 when 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 Sleeping until the deadline. A delivery whose payload matches the JSON schema of the type wakes it, and the handler receives Some(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 rejected in the response. The signal is still stored.
  • Idempotency: a second delivery with the same idempotency_id returns duplicate: true and 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

ChannelHow
RESTPOST /api/v1/signals with {"name", "key", "payload", "idempotency_id"}; admin JWT or API key with the signals_send scope
RESTGET /api/v1/signals?name=&key= lists received signals (runs_read for an API key)
CLIironflow signal send ci.pipeline_finished --key <sha> --payload '{"status":"success"}'
CLIironflow signal list --name ci.pipeline_finished
MCPsend_signal, list_signals
Rustengine.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.

TypeField attributeField typeAnswer
noul#[noul("..")], optionally if_true = "..", if_false = ".."f64probability of “yes” in [0, 1]
choice#[choice("..")]an enum deriving DecisionChoicethe option picked
score#[score("..", levels = ["..", ".."])]f64probability-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", &notifier).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:

  1. Add the dependency to your workflow crate
  2. Build a client from your workflow’s OperationContext (via from_context())
  3. 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

CrateDescriptioncargo add
ironflow-ops-commonShared HTTP client and helpers for ops crates (not used directly in workflows)cargo add ironflow-ops-common
ironflow-ops-dockerDocker containers, images, networks, volumes via bollardcargo add ironflow-ops-docker
ironflow-ops-gitGit operations (commit, branch, merge, diff, …) via git2cargo add ironflow-ops-git
ironflow-ops-gitlabGitLab API v4 (issues, MRs, pipelines, …) via the gitlab cratecargo add ironflow-ops-gitlab
ironflow-ops-grafanaGrafana API (dashboards, alerting, data sources, …)cargo add ironflow-ops-grafana
ironflow-ops-helmHelm CLI wrapper (install, upgrade, rollback, charts, repos)cargo add ironflow-ops-helm
ironflow-ops-k8sKubernetes typed API via kube + k8s-openapi, plus run-to-completion ops (PodRun, JobRun, ApplyConfigMap, ApplySecret)cargo add ironflow-ops-k8s
ironflow-ops-lokiGrafana Loki (log queries, ingest, rules, labels)cargo add ironflow-ops-loki
ironflow-ops-mimirGrafana Mimir (PromQL queries, remote write, rules, cardinality)cargo add ironflow-ops-mimir
ironflow-ops-postgresPostgreSQL queries and admin via sqlxcargo add ironflow-ops-postgres
ironflow-ops-s3AWS S3 objects, buckets, presigned URLs via aws-sdk-s3cargo add ironflow-ops-s3
ironflow-ops-slackSlack API (chat, conversations, files, users) via slack-morphismcargo add ironflow-ops-slack
ironflow-ops-tempoGrafana 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:

  1. Register secrets in your workflow’s secret store
  2. Call from_context(&ctx) which resolves them automatically
  3. 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; each output reads 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_fast is true (the second argument), the remaining steps are cancelled when one fails
  • If fail_fast is false, 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::StepConfig before 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 Operation is 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:

FlagMeaning
--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-estimatesSkip the duration estimate query against run history.
--jsonEmit 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:

StateMeaning
evaluatedDeclared with ctx.when and resolved against the input.
skippedThe handler called ctx.skip(name, reason) on this branch.
unevaluableDeclared 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

TransportProviderUse case
LocalClaudeCodeProviderDevelopment, simple setups
DockerDockerProviderIsolated containers on the same host
SSHSshProviderRemote machines
KubernetesK8sProviderEphemeral 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.

PathContent
/home/claudeHOME of uid 10001, an emptyDir in the sandbox
/tmpTMPDIR, an emptyDir in the sandbox
/etc/claude-codemanaged-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

DefaultRelaxation
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 textnone: 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 receives ANTHROPIC_BASE_URL (the proxy URL), ANTHROPIC_AUTH_TOKEN (the opaque token) and CLAUDE_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/ to api.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_subscription account kind once auth_proxy is set), else CLAUDE_CODE_OAUTH_TOKEN, then ANTHROPIC_API_KEY, from the worker environment. The worker also needs IRONFLOW_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_proxy set, a Claude credential configured for the pod (oauth_token_from_secret, oauth_credentials, oauth_credentials_from_secret, a step env_from_secret("CLAUDE_CODE_OAUTH_TOKEN", ..), any plain value starting with sk-ant) fails the step with an error naming the variable.
  • Registry. In memory by default (one replica). Set IRONFLOW_AUTH_PROXY_DATABASE_URL and 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) of cilium-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 PodRun and Jobs of JobRun from ironflow-ops-k8s (.expiry_margin(d), 60s by default). A caller cannot set managed-by or component: pod_label, PodRun::label and JobRun::label panic, 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, cargo and npm installs). Keep gVisor for the steps that handle untrusted code and leave trusted build steps on runc.

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:

ProductionUnder TestEngine
API servernothing to start; the run executes inline
Background workernothing to start; run() returns once the run is finished
PostgresInMemoryStore
sh -c <command>a closure
An HTTP requesta closure
The Claude CLIa closure, or a recorded fixture
A human clicking Approvean 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

MethodWhat 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:

MethodWhat 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.

AccessorReturns
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 MockShellOutput with a non-zero exit_code is an error, like a real non-zero exit. Use MockShellOutput::failed(1, "boom"). When the step sets exit_code_as_output(), the mock completes with that exit code as its output.
  • A MockHttpResponse with a non-2xx status is a normal output, like a real 500 response. Return Err(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-double Operation to the handler.
  • ctx.delay(...) is not intercepted: a non-zero delay still suspends the run with RunStatus::Sleeping. It resumes once RunWaker::tick runs after its scheduled_at.
  • ctx.wait_for_signal(...) is resolved with with_mock_signal(|step, name, key| ..), returning SignalOutcome::Received(json!(..)) or SignalOutcome::TimedOut. Without it, the run ends in RunStatus::Sleeping until a signal is delivered.
  • ctx.decision(...) needs a real DecisionProvider, wired with with_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

CrateRole
ironflow-coreShell execution, agent providers, cost tracking
ironflow-engineWorkflow handler trait, context, step orchestration
ironflow-apiREST API (axum), routes, SSE events, dashboard serving
ironflow-workerBackground worker that polls and executes runs
ironflow-storeStorage trait + PostgreSQL and in-memory backends
ironflow-authJWT authentication, password hashing, API keys
ironflow-runtimeDaemon features: webhooks, trigger sources
ironflow-artifactsBlob storage for step-produced files
ironflow-templatesFetch and install workflow templates from Git
ironflow-sdkType-safe Rust client (types generated from OpenAPI)
ironflow-cliCommand-line interface (clap v4)
ironflow-mcpModel Context Protocol server
ironflow-typesShared API envelope types

Request flow

  1. A client (dashboard, CLI, SDK, or webhook) sends a request to the API
  2. The API validates authentication, creates a Run in the Store, and returns it
  3. A Worker polls the API, acquires a lease on the Run, and executes the handler
  4. The handler calls steps (ctx.shell(), ctx.agent(), etc.), each persisted as they complete
  5. Events are published via SSE for real-time updates
  6. 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.