A lightweight, LangGraph-inspired stateful workflow engine for Rust.
graphflow enables you to coordinate multi-step workflows, autonomous AI agent loops, and state machines with ease. Model your processes as directed graphs with nodes (actions), static edges (transitions), and conditional edges (dynamic routing and decision making).
- Features
- How It Works
- Installation
- Quickstart Tutorial
- Core Concepts
- Examples
- API Reference
- Graph Validation Rules
- Running Tests & Examples
- License
- β‘ Zero Heavy Dependencies: Pure Rust with standard library hash maps and function pointers.
- π Type-Safe State Machine: Graph execution is generic over your custom state struct
T. - π Dynamic Branching & Decision Routing: Prioritizes conditional edges to route execution based on runtime state inspection.
- π Cycles & Loops: Seamlessly supports looping workflows for agent retries, human-in-the-loop flows, or iterative optimization.
- π‘οΈ Pre-Execution Graph Validation: Verifies entry/finish nodes and edge integrity during
.compile()before any execution begins. - π¦ Comprehensive Custom Errors: Strongly-typed
GraphError,PregelError,ChannelError, and unifiedGraphflowErrorreplacing basic/string errors. - π Built-in Styled Logging: Powered by the
logcrate withinit_logger()providing colored levels, timestamps, and module targets. - π οΈ Fluent Builder API: Intuitive method chaining to construct workflows cleanly.
A graphflow graph consists of:
- Nodes: Functions that inspect state and mutate it safely through
StateContext<'_, T>. - Edges: Deterministic paths connecting one node to the next.
- Conditional Edges: Decision functions that inspect
&Tand dynamically choose the next target node. - Lifecycle Hooks: Observers that track execution events (
on_agent_start,on_agent_state_change,on_node_start,on_node_end,on_conditional_node_start,on_conditional_node_end,on_agent_end).
flowchart LR
Start([Entry Point: Start]) --> A[Node A]
A --> B{Conditional Edge}
B -- "Needs more work" --> A
B -- "Ready" --> C[Node C]
C --> Finish([Finish Point])
During execution:
- Execution begins at the configured
entry_point. - At each node, the corresponding function executes and mutates the shared state.
- If the current node is the
finish_point, execution terminates successfully. - Otherwise, if a conditional edge exists for the current node, it is evaluated first.
- If no conditional edge exists, the standard static edge is followed.
Add graphflow to your Cargo.toml:
[dependencies]
graphflow = "0.3.0"Or add it from your local workspace / git repository:
[dependencies]
graphflow = { git = "https://github.com/vesal-j/graphflow" }Here is a complete, minimal working example in 5 steps:
use graphflow::{init_logger, Graph, GraphError, StateContext};
// Step 1: Define your state
struct WorkflowState {
pub message: String,
pub step_count: usize,
}
// Step 2: Define your node functions using StateContext and custom GraphError
fn step_one(ctx: &mut StateContext<'_, WorkflowState>) -> Result<(), GraphError> {
ctx.update(|state| {
state.message.push_str("Hello");
state.step_count += 1;
});
Ok(())
}
fn step_two(ctx: &mut StateContext<'_, WorkflowState>) -> Result<(), GraphError> {
ctx.update(|state| {
state.message.push_str(" World!");
state.step_count += 1;
});
Ok(())
}
fn main() -> Result<(), Box<dyn std::error::Error>> {
// Optional: Initialize the styled logging system
init_logger();
// Step 3: Build the graph
let mut graph: Graph<WorkflowState> = Graph::new();
graph
.add_node("first".to_string(), step_one)
.add_node("second".to_string(), step_two)
.add_edge("first".to_string(), "second".to_string())
.set_entry_point("first".to_string())
.set_finish_point("second".to_string());
// Step 4: Compile and validate the graph (returns Result<CompiledGraph<T>, GraphError>)
let compiled = graph.compile()?;
// Step 5: Execute the graph with initial state
let state = WorkflowState {
message: String::new(),
step_count: 0,
};
let final_state = compiled.invoke(state)?;
println!("Final message: {}", final_state.message);
Ok(())
}The state represents the central blackboard or context of your pipeline. Nodes inspect and mutate state through StateContext<'_, T>, ensuring that state changes are controlled and tracked by the hook system.
struct AgentState {
user_prompt: String,
tool_results: Vec<String>,
iteration: usize,
}Nodes perform discrete units of computation. Instead of receiving a raw mutable reference, nodes receive a &mut StateContext<'_, T>:
pub type NodeFunction<T> = fn(&mut StateContext<'_, T>) -> Result<(), GraphError>;- Read state: Access fields directly using
ctx.field(viaDeref) orctx.get(). - Mutate state: Call
ctx.update(|state| { ... })orctx.set(new_state). Every mutation automatically triggers theon_agent_state_changelifecycle hook. Read-only nodes will not trigger unnecessary state change events.
Example:
fn execute_tool(ctx: &mut StateContext<'_, AgentState>) -> Result<(), GraphError> {
log::info!("Current iteration: {}", ctx.iteration); // Read via Deref
ctx.update(|state| {
state.tool_results.push("Tool output".to_string());
}); // Automatically triggers on_agent_state_change
Ok(())
}Static edges create an unconditional transition from node from to node to:
graph.add_edge("fetch_data".to_string(), "process_data".to_string());Conditional edges evaluate the current state and return the name of the next node:
pub type BranchFunction<T> = fn(&T) -> String;Note: If a node has both a static edge and a conditional edge, the conditional edge takes precedence.
graph.add_conditional_edge("evaluate".to_string(), |state| {
if state.iteration < 3 {
"retry".to_string()
} else {
"finish".to_string()
}
});Every graph must specify where execution starts and where it ends:
graph.set_entry_point("start_node".to_string());
graph.set_finish_point("final_node".to_string());Calling .compile() verifies graph integrity before executing:
let compiled = graph.compile()?;If any referenced node is missing or entry/finish points are undefined, compile returns a descriptive typed Err(GraphError) (e.g. GraphError::MissingEntryPoint, GraphError::EdgeSourceNotFound).
In this example, an integer state is incremented and doubled in a loop until it reaches a threshold:
use graphflow::Graph;
struct CounterState {
count: usize,
}
fn main() -> Result<(), Box<dyn std::error::Error>> {
let mut graph: Graph<CounterState> = Graph::new();
graph
.add_node("increment".to_string(), |ctx| {
ctx.update(|state| state.count += 1);
println!("Count incremented to {}", ctx.count);
Ok(())
})
.add_node("double".to_string(), |ctx| {
ctx.update(|state| state.count *= 2);
println!("Count doubled to {}", ctx.count);
Ok(())
})
.add_node("finish".to_string(), |ctx| {
println!("Done! Final count: {}", ctx.count);
Ok(())
})
.add_edge("increment".to_string(), "double".to_string())
.add_conditional_edge("double".to_string(), |state| {
if state.count < 10 {
"increment".to_string() // loop back
} else {
"finish".to_string() // exit loop
}
})
.set_entry_point("increment".to_string())
.set_finish_point("finish".to_string());
let compiled = graph.compile().map_err(|e| format!("Validation error: {e}"))?;
compiled.invoke(CounterState { count: 1 })?;
Ok(())
}Run this example directly from the repository:
cargo run --example counter_loopModel a classic AI agent loop:
- Plan: Formulate plan or analyze prompt.
- Search / Tool: Gather external information.
- Route: Decide whether to call tools again or synthesize the answer.
- Synthesize: Produce final response.
use graphflow::{init_logger, Graph, GraphError, StateContext};
struct AgentState {
query: String,
iterations: usize,
max_iterations: usize,
has_sufficient_context: bool,
response: Option<String>,
}
fn plan(ctx: &mut StateContext<'_, AgentState>) -> Result<(), GraphError> {
log::info!("[Plan] Planning query: '{}'", ctx.query);
Ok(())
}
fn search_tools(ctx: &mut StateContext<'_, AgentState>) -> Result<(), GraphError> {
ctx.update(|state| {
state.iterations += 1;
log::info!("[Tool] Running search (iteration {})...", state.iterations);
if state.iterations >= 2 {
state.has_sufficient_context = true;
}
});
Ok(())
}
fn route_decision(state: &AgentState) -> String {
if state.has_sufficient_context || state.iterations >= state.max_iterations {
"synthesize".to_string()
} else {
"search".to_string()
}
}
fn synthesize(ctx: &mut StateContext<'_, AgentState>) -> Result<(), GraphError> {
log::info!("[Synthesize] Synthesizing final response...");
ctx.update(|state| {
state.response = Some(format!("Answer for: {}", state.query));
});
Ok(())
}
fn main() -> Result<(), Box<dyn std::error::Error>> {
init_logger();
let mut graph: Graph<AgentState> = Graph::new();
graph
.add_node("plan".to_string(), plan)
.add_node("search".to_string(), search_tools)
.add_node("synthesize".to_string(), synthesize)
.add_edge("plan".to_string(), "search".to_string())
.add_conditional_edge("search".to_string(), route_decision)
.set_entry_point("plan".to_string())
.set_finish_point("synthesize".to_string());
let compiled = graph.compile()?;
let state = AgentState {
query: "What is the capital of Rustland?".to_string(),
iterations: 0,
max_iterations: 3,
has_sufficient_context: false,
response: None,
};
compiled.invoke(state)?;
Ok(())
}Run this example:
cargo run --example llm_agent_workflowThe builder struct used to declare workflow topology.
| Method | Signature | Description |
|---|---|---|
new() |
pub fn new() -> Self |
Creates an empty graph builder. |
default() |
fn default() -> Self |
Equivalent to Graph::new(). |
add_node |
&mut Self, name: String, function: NodeFunction<T> -> &mut Self |
Registers a named node function. |
add_edge |
&mut Self, from: String, to: String -> &mut Self |
Adds a directional edge between two nodes. |
add_conditional_edge |
&mut Self, from: String, branch: BranchFunction<T> -> &mut Self |
Adds a dynamic branching edge evaluated via BranchFunction<T>. |
set_entry_point |
&mut Self, from: String -> &mut Self |
Sets the start node name. |
set_finish_point |
&mut Self, from: String -> &mut Self |
Sets the termination node name. |
compile |
&mut Self -> Result<CompiledGraph<T>, GraphError> |
Validates graph structure and returns an executable CompiledGraph. |
The validated, immutable graph ready for execution.
| Method | Signature | Description |
|---|---|---|
invoke |
&self, state: T -> Result<T, GraphError> |
Traverses the graph, executing nodes and routing decisions, and returns the final state. |
graphflow uses strongly-typed custom error enums instead of basic Rust strings or placeholder errors:
GraphError: Errors during compilation and execution ofGraph<T>/CompiledGraph<T>:MissingEntryPoint/MissingFinishPoint: Entry or finish point was not configured.EntryPointNotFound(String)/FinishPointNotFound(String): Specified node does not exist.EdgeSourceNotFound(String)/EdgeTargetNotFound(String): Edge endpoints are missing.ConditionalEdgeSourceNotFound(String): Branching source node does not exist.NodeNotFound(String): Referenced node not present in graph during traversal.MissingEdge(String): Reached a dead-end non-finish node with no outgoing edge.ExecutionError(String): Node function error during execution.
PregelError: Errors in the Bulk Synchronous Parallel execution engine:ChannelNotFound(String)/MissingChannel(String)/NodeNotFound(String).TypeMismatch(String): Data type mismatch on channel write.MaxStepsExceeded(usize): Cycle did not terminate withinmax_steps.ExecutionError(String): Worker thread or node execution failed.Channel(ChannelError): Underlying channel operation failed.
ChannelError: Type downcast mismatches and channel updates.GraphflowError: Unified top-level error enum implementingstd::error::Errorand source error chaining.
Initialize styled, timestamped, colored logging powered by the log crate:
use graphflow::{init_logger, init_logger_with_level};
use log::LevelFilter;
// Initialize with default INFO level (honors RUST_LOG env var)
init_logger();
// Or specify a default filter level
init_logger_with_level(LevelFilter::Debug);graphflow provides an event-driven hook system to observe the complete graph lifecycle without coupling business logic to external event dispatchers. All hooks receive read-only references to state (&T):
| Hook Event | Trigger Point | Arguments |
|---|---|---|
on_agent_start |
Graph execution begins at entry point | entry_point: &str, state: &T |
on_node_start |
Immediately before a node begins | node_name: &str, state: &T |
on_agent_state_change |
Triggered whenever state is mutated via ctx.update or ctx.set |
node_name: &str, state: &T |
on_node_end |
Immediately after a node finishes | node_name: &str, state: &T, result: &Result<(), GraphError> |
on_conditional_node_start |
Before evaluating a conditional edge | from_node: &str, state: &T |
on_conditional_node_end |
After dynamic routing decision is made | from_node: &str, next_node: &str, state: &T |
on_agent_end |
Graph execution completes at finish point | finish_point: &str, state: &T |
graph
.on_agent_start(|entry, state| println!("Agent starting at '{entry}'"))
.on_agent_state_change(|node, state| println!("State changed in '{node}'"))
.on_node_start(|node, state| println!("Node '{node}' starting"))
.on_node_end(|node, state, result| println!("Node '{node}' finished: {:?}", result.is_ok()))
.on_agent_end(|finish, state| println!("Agent finished at '{finish}'"));use graphflow::{GraphHook, GraphError};
struct AuditHook;
impl<T: std::fmt::Debug> GraphHook<T> for AuditHook {
fn on_agent_start(&self, entry: &str, state: &T) {
println!("[Audit] Start: {entry}, state: {state:?}");
}
fn on_agent_state_change(&self, node: &str, state: &T) {
println!("[Audit] State modified by {node}: {state:?}");
}
}
graph.add_hook(AuditHook);When .compile() is invoked, graphflow ensures:
- Entry Point Configured: An entry point must be provided (
"Entry point is not set"). - Finish Point Configured: A finish point must be provided (
"Finish point is not set"). - Valid Entry Node: Entry point (or all comma-separated entry points) must correspond to added nodes (
"Entry point '<name>' does not exist"). - Valid Finish Node: Finish point must correspond to an added node (
"Finish point '<name>' does not exist"). - Edge Integrity: For every edge
from -> to:frommust exist ("Edge source '<name>' does not exist").- all targets in
tomust exist ("Edge target '<name>' does not exist").
- Conditional Edge Integrity: For every conditional edge
from:frommust exist ("Conditional edge source '<name>' does not exist").
To run the complete test suite (39 tests):
cargo testTo run all bundled examples:
# Run counter loop example
cargo run --example counter_loop
# Run LLM agent workflow example
cargo run --example llm_agent_workflow
# Run Pregel parallel Fork-Join example
cargo run --example parallel_pregelThis project is licensed under the MIT License.