Build agents, workflows, and workers in Rust with Conductor.
See the agent quickstart or workflow quickstart. Conductor Skills helps coding agents work with Conductor.
Use a hosted or local Conductor server for the examples below.
In Orkes Developer Edition, create an application and access key, then set:
export CONDUCTOR_SERVER_URL=https://developer.orkescloud.com/api
export CONDUCTOR_AUTH_KEY=<your-key-id>
export CONDUCTOR_AUTH_SECRET=<your-key-secret>For another remote cluster, use its /api URL and credentials.
npm install -g @conductor-oss/conductor-cli
conductor server start
conductor server status
export CONDUCTOR_SERVER_URL=http://localhost:8080/apidocker compose -f scripts/docker-compose-oss.yaml up -d
export CONDUCTOR_SERVER_URL=http://localhost:8080/apiThe Compose server UI is at http://localhost:8080.
Requires Rust 1.85 or newer and a reachable Conductor server. See CI for tested server versions.
Add the following to your Cargo.toml:
[dependencies]
# The crate is published as `conductor-sdk`; rename it to `conductor` so
# `use conductor::...` works in your code.
conductor = { version = "0.2.0", package = "conductor-sdk" }
tokio = { version = "1", features = ["full"] }For the #[worker] macro (similar to Python's @worker_task decorator):
[dependencies]
conductor = { version = "0.2.0", package = "conductor-sdk", features = ["macros"] }
conductor-macros = "0.2.0"
tokio = { version = "1", features = ["full"] }For agents, enable agents:
conductor = { version = "0.2.0", package = "conductor-sdk", features = ["agents"] }
tokio = { version = "1", features = ["full"] }The agents feature includes worker macros.
Configure a provider and model on your Conductor server. Create a project and add the agents dependencies above to its Cargo.toml:
cargo new conductor-agent-demo
cd conductor-agent-demoPut this in src/main.rs:
use conductor::agents::{AgentDef, AgentRuntime};
use conductor::configuration::Configuration;
use conductor::error::{ConductorError, Result};
#[tokio::main]
async fn main() -> Result<()> {
let model = std::env::args().nth(1).ok_or_else(|| {
ConductorError::agent("Pass a configured provider/model as the first argument")
})?;
if model.trim().is_empty() {
return Err(ConductorError::agent("Model argument cannot be empty"));
}
let config = Configuration::from_env();
let runtime = AgentRuntime::new(config)?;
let agent = AgentDef::new("greeter")?
.with_model(model)
.with_instructions("You are a friendly assistant. Keep responses brief.");
let result = runtime
.run(
&agent,
"Say hello and tell me a fun fact about Python.".into(),
)
.await?;
println!("status: {}", result.status);
println!("output: {}", result.output);
println!("execution_id: {}", result.execution_id);
if !result.is_success() {
return Err(ConductorError::agent(format!(
"Agent execution ended with status {}",
result.status
)));
}
Ok(())
}Run with a model configured on your server:
cargo run -- openai/gpt-4o-miniThe checked-in example uses the same code.
Expect status: COMPLETED, model output, and an execution ID.
See the agent guide and examples for more.
Step 1: Create a workflow
Workflows are definitions that reference task types (e.g. a SIMPLE task called greet). We'll build a workflow called
greetings that runs one task and returns its output.
use conductor::models::{WorkflowDef, WorkflowTask};
fn greetings_workflow() -> WorkflowDef {
WorkflowDef::new("greetings")
.with_version(1)
.with_task(
WorkflowTask::simple("greet", "greet_ref")
.with_input_param("name", "${workflow.input.name}")
)
.with_output_param("result", "${greet_ref.output.result}")
}Step 2: Write worker
Workers are Rust functions decorated with #[worker] that poll Conductor for tasks and execute them.
use conductor_macros::worker;
#[worker(name = "greet")]
async fn greet(name: String) -> String {
format!("Hello {}", name)
}Step 3: Run your first workflow app
Create a Cargo project, add the macro dependencies above to its Cargo.toml, then put the following in src/main.rs:
cargo new greetings
cd greetingsuse conductor::{
client::ConductorClient,
configuration::Configuration,
models::{StartWorkflowRequest, WorkflowDef, WorkflowTask},
worker::TaskHandler,
};
use conductor_macros::worker;
// A worker is any Rust function with the #[worker] macro.
#[worker(name = "greet")]
async fn greet(name: String) -> String {
format!("Hello {}", name)
}
fn greetings_workflow() -> WorkflowDef {
WorkflowDef::new("greetings")
.with_version(1)
.with_task(
WorkflowTask::simple("greet", "greet_ref")
.with_input_param("name", "${workflow.input.name}")
)
.with_output_param("result", "${greet_ref.output.result}")
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
// Configure the SDK (reads CONDUCTOR_SERVER_URL / CONDUCTOR_AUTH_* from env).
let config = Configuration::default();
let client = ConductorClient::new(config.clone())?;
// Register the workflow
let workflow = greetings_workflow();
client.metadata_client()
.register_or_update_workflow_def(&workflow, true)
.await?;
// Start polling for tasks
let mut task_handler = TaskHandler::new(config.clone())?;
task_handler.add_worker(greet_worker());
task_handler.start().await?;
// Run the workflow and get the result
let run = client.workflow_client()
.execute_workflow(
&StartWorkflowRequest::new("greetings")
.with_version(1)
.with_input_value("name", "Conductor"),
std::time::Duration::from_secs(10),
)
.await?;
println!("result: {:?}", run.output.get("result"));
println!("execution: {}/execution/{}", config.ui_host, run.workflow_id);
task_handler.stop().await?;
Ok(())
}Run it:
cargo runExport your authentication credentials as well:
export CONDUCTOR_SERVER_URL="https://your-cluster.orkesconductor.io/api" # If using Orkes Conductor that requires auth key/secret export CONDUCTOR_AUTH_KEY="your-key" export CONDUCTOR_AUTH_SECRET="your-secret"See Choose your Conductor server for setup details.
That's it -- you just defined a worker, built a workflow, and executed it. Open the Conductor UI (default: http://localhost:8080) to see the execution.
The example includes sync + async workers, metrics, and long-running tasks.
See examples/worker_example.rs
Workers are Rust functions that execute Conductor tasks. Use the #[worker] macro or FnWorker to:
- register it as a worker (auto-discovered by
TaskHandler) - use it as a workflow task (call it with
task_ref_name=...)
Workers can also serve as agent tools. See Conductor agents.
use conductor_macros::worker;
#[worker(name = "greet")]
async fn greet(name: String) -> String {
format!("Hello {}", name)
}Using FnWorker (closure-based):
use conductor::worker::{FnWorker, WorkerOutput};
let greetings_worker = FnWorker::new("greetings", |task| async move {
let name = task.get_input_string("name").unwrap_or_default();
Ok(WorkerOutput::completed_with_result(format!("Hello, {}", name)))
})
.with_thread_count(10)
.with_poll_interval_millis(100);Start workers with TaskHandler:
use conductor::{
configuration::Configuration,
worker::TaskHandler,
};
let config = Configuration::default();
let mut task_handler = TaskHandler::new(config)?;
task_handler.add_worker(greet_worker());
task_handler.start().await?;
// Wait for shutdown signal
tokio::signal::ctrl_c().await?;
task_handler.stop().await?;Worker Configuration
Workers support hierarchical environment variable configuration — global settings that can be overridden per worker:
# Global (all workers)
export CONDUCTOR_WORKER_ALL_POLL_INTERVAL_MILLIS=250
export CONDUCTOR_WORKER_ALL_THREAD_COUNT=20
export CONDUCTOR_WORKER_ALL_DOMAIN=production
# Per-worker override
export CONDUCTOR_WORKER_GREETINGS_THREAD_COUNT=50See WORKER_CONFIGURATION.md for all options.
Enable Prometheus metrics:
use conductor::metrics::MetricsSettings;
use conductor::worker::TaskHandler;
let mut task_handler = TaskHandler::new(config)?;
task_handler.enable_metrics(
MetricsSettings::new()
.with_http_port(9090)
);
task_handler.start().await?;
// Metrics at http://localhost:9090/metricsSee METRICS.md for details.
Learn more:
- Worker Guide — All worker patterns (function, closure, macro, async)
- Worker Configuration — Environment variable configuration system
Define workflows in Rust using the builder pattern to chain tasks:
use conductor::{
client::ConductorClient,
configuration::Configuration,
models::{WorkflowDef, WorkflowTask},
};
let config = Configuration::default();
let client = ConductorClient::new(config)?;
let metadata_client = client.metadata_client();
let workflow = WorkflowDef::new("greetings")
.with_version(1)
.with_task(
WorkflowTask::simple("greet", "greet_ref")
.with_input_param("name", "${workflow.input.name}")
)
.with_output_param("result", "${greet_ref.output.result}");
// Registering is required if you want to start/execute by name+version
metadata_client.register_or_update_workflow_def(&workflow, true).await?;Execute workflows:
use conductor::models::StartWorkflowRequest;
use std::time::Duration;
// Asynchronous (returns workflow ID immediately)
let request = StartWorkflowRequest::new("greetings")
.with_version(1)
.with_input_value("name", "Orkes");
let workflow_id = workflow_client.start_workflow(&request).await?;
// Synchronous (waits for completion)
let run = workflow_client
.execute_workflow(&request, Duration::from_secs(10))
.await?;
println!("{:?}", run.output);Manage running workflows and send signals:
workflow_client.pause_workflow(&workflow_id).await?;
workflow_client.resume_workflow(&workflow_id).await?;
workflow_client.terminate_workflow(&workflow_id, Some("no longer needed"), false).await?;
workflow_client.retry_workflow(&workflow_id, false).await?;
workflow_client.restart_workflow(&workflow_id, false).await?;Learn more:
- Workflow Management — Start, pause, resume, terminate, retry, search
- Workflow Testing — Example using mock task outputs
- Metadata Management — Task & workflow definitions
- Worker stops polling:
TaskHandlermonitors workers. Usetask_handler.is_healthy()for health checks. - Connection issues: Verify
CONDUCTOR_SERVER_URLis correct and server is running. - Authentication failures: For Orkes Conductor, ensure
CONDUCTOR_AUTH_KEYandCONDUCTOR_AUTH_SECRETare valid. - Agent cannot call a model: Check the provider and model configured on the Conductor server.
The agents feature supports local tools, human approval, guardrails, and multi-agent runs.
Configure a model on your Conductor server, then try:
| Example | Shows |
|---|---|
| Simple tools | Rust function tools |
| Human approval | Terminal approval for a tool call |
| Parallel agents | Multi-agent execution |
cargo run --example agent_demo_02a_simple_tools --features agents -- openai/gpt-4o-miniSee the agent guide and full example list. For LLM and RAG workflows built with the core SDK, see llm_chat_example.rs and rag_workflow.rs.
See the examples directory for the full catalog. Key examples:
| Example | Description | Run |
|---|---|---|
| worker_example.rs | End-to-end: sync + async workers, metrics | cargo run --example worker_example |
| hello_world.rs | Minimal hello world | cargo run --example hello_world |
| dynamic_workflow.rs | Build workflows programmatically | cargo run --example dynamic_workflow |
| llm_chat_example.rs | AI multi-turn chat | cargo run --example llm_chat_example |
| rag_workflow.rs | RAG pipeline | cargo run --example rag_workflow |
| task_context_example.rs | Long-running tasks with TaskContext | cargo run --example task_context_example |
| workflow_ops.rs | Pause, resume, terminate workflows | cargo run --example workflow_ops |
| test_workflows.rs | Unit testing workflows | cargo run --example test_workflows |
| kitchensink.rs | All task types (HTTP, JS, JQ, Switch) | cargo run --example kitchensink |
End-to-end examples covering all APIs for each domain:
| Example | APIs | Run |
|---|---|---|
| authorization_example.rs | Authorization APIs | cargo run --example authorization_example |
| metadata_journey.rs | Metadata APIs | cargo run --example metadata_journey |
| schedule_journey.rs | Schedule APIs | cargo run --example schedule_journey |
| prompt_journey.rs | Prompt APIs | cargo run --example prompt_journey |
| Document | Description |
|---|---|
| Worker Guide | All worker patterns (function, closure, macro, async) |
| Agent Guide | Durable agents, tools, and runtime modes |
| Worker Configuration | Hierarchical environment variable configuration |
| Workflow Management | Start, pause, resume, terminate, retry, search |
| Workflow Testing | Example using mock task outputs |
| Task Management | Task operations |
| Metadata | Task & workflow definitions |
| Authorization | Users, groups, applications, permissions |
| Schedules | Workflow scheduling |
| Secrets | Secret storage |
| Prompts | AI/LLM prompt templates |
| Integrations | AI/LLM provider integrations |
| Metrics | Prometheus metrics collection |
- Open an issue (SDK) for SDK bugs, questions, and feature requests
- Open an issue (Conductor server) for Conductor OSS server issues
- Join the Conductor Slack for community discussion and help
- Orkes Community Forum for Q&A
- Conductor OSS contribution guide for contributing upstream
- Conductor Code of Conduct and security policy for community and private vulnerability reporting
Is this the same as Netflix Conductor?
Yes. Conductor OSS is the continuation of the original Netflix Conductor repository after Netflix contributed the project to the open-source foundation.
Is this project actively maintained?
Yes. Orkes is the primary maintainer and offers an enterprise SaaS platform for Conductor across all major cloud providers.
Can Conductor scale to handle my workload?
Conductor was built at Netflix to handle massive scale and has been battle-tested in production environments processing millions of workflows. It scales horizontally to meet virtually any demand.
Does Conductor support durable code execution?
Yes. Conductor ensures workflows complete reliably even in the face of infrastructure failures, process crashes, or network issues.
Are workflows always asynchronous?
No. While Conductor excels at asynchronous orchestration, it also supports synchronous workflow execution when immediate results are required.
Do I need to use a Conductor-specific framework?
No. Conductor is language and framework agnostic. Use your preferred language and framework -- the SDKs provide native integration for Python, Java, JavaScript, Go, C#, Rust, and more.
Can I mix workers written in different languages?
Yes. A single workflow can have workers written in Rust, Python, Java, Go, or any other supported language. Workers communicate through the Conductor server, not directly with each other.
What Rust versions are supported?
Rust 1.85 and above (2021 edition), as declared in Cargo.toml.
Should I use async fn or regular fn for my workers?
Use async fn for I/O-bound tasks (API calls, database queries) — the SDK uses async runtime for high concurrency with low overhead. Use regular functions for CPU-bound or blocking work. The SDK handles both patterns efficiently.
How do I run workers in production?
Workers are standard Rust applications. Deploy them as you would any Rust application -- in containers, VMs, or bare metal. Workers poll the Conductor server for tasks, so no inbound ports need to be opened.
How do I test workflows with mock task outputs?
Conductor's POST /api/workflow/test endpoint evaluates workflows with mock task outputs. See the workflow testing example. This endpoint still requires a running Conductor server.
Apache 2.0