Skip to content

Latest commit

Β 

History

17 Commits

Folders and files

NameName
Last commit message
Last commit date
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 
Β 

Repository files navigation

graphflow πŸ¦€β›“οΈ

Crates.io License: MIT Rust: 2024 Tests

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


πŸ“‘ Table of Contents


✨ Features

  • ⚑ 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 unified GraphflowError replacing basic/string errors.
  • πŸ“œ Built-in Styled Logging: Powered by the log crate with init_logger() providing colored levels, timestamps, and module targets.
  • πŸ› οΈ Fluent Builder API: Intuitive method chaining to construct workflows cleanly.

🧠 How It Works

A graphflow graph consists of:

  1. Nodes: Functions that inspect state and mutate it safely through StateContext<'_, T>.
  2. Edges: Deterministic paths connecting one node to the next.
  3. Conditional Edges: Decision functions that inspect &T and dynamically choose the next target node.
  4. 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])
Loading

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.

πŸ“¦ Installation

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" }

πŸš€ Quickstart Tutorial

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(())
}

🧩 Core Concepts

1. Shared State (T)

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,
}

2. Nodes (NodeFunction<T>) & StateContext

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 (via Deref) or ctx.get().
  • Mutate state: Call ctx.update(|state| { ... }) or ctx.set(new_state). Every mutation automatically triggers the on_agent_state_change lifecycle 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(())
}

3. Static Edges

Static edges create an unconditional transition from node from to node to:

graph.add_edge("fetch_data".to_string(), "process_data".to_string());

4. Conditional Edges (BranchFunction<T>)

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()
    }
});

5. Entry and Finish Points

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());

6. Compilation & Validation

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


πŸ’‘ Examples

Example 1: Cyclic Loop with Dynamic Routing

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_loop

Example 2: Autonomous LLM Agent Loop

Model a classic AI agent loop:

  1. Plan: Formulate plan or analyze prompt.
  2. Search / Tool: Gather external information.
  3. Route: Decide whether to call tools again or synthesize the answer.
  4. 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_workflow

πŸ“– API Reference

Graph<T>

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

CompiledGraph<T>

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.

🚦 Error Handling & Custom Types

graphflow uses strongly-typed custom error enums instead of basic Rust strings or placeholder errors:

  • GraphError: Errors during compilation and execution of Graph<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 within max_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 implementing std::error::Error and source error chaining.

πŸ“œ Logging System

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);

πŸͺ Lifecycle Hook System

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

Using Closure Hooks:

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}'"));

Using Custom Struct Hook (GraphHook<T>):

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);

πŸ›‘οΈ Graph Validation Rules

When .compile() is invoked, graphflow ensures:

  1. Entry Point Configured: An entry point must be provided ("Entry point is not set").
  2. Finish Point Configured: A finish point must be provided ("Finish point is not set").
  3. Valid Entry Node: Entry point (or all comma-separated entry points) must correspond to added nodes ("Entry point '<name>' does not exist").
  4. Valid Finish Node: Finish point must correspond to an added node ("Finish point '<name>' does not exist").
  5. Edge Integrity: For every edge from -> to:
    • from must exist ("Edge source '<name>' does not exist").
    • all targets in to must exist ("Edge target '<name>' does not exist").
  6. Conditional Edge Integrity: For every conditional edge from:
    • from must exist ("Conditional edge source '<name>' does not exist").

πŸ§ͺ Running Tests & Examples

To run the complete test suite (39 tests):

cargo test

To 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_pregel

πŸ“œ License

This project is licensed under the MIT License.

About

A lightweight, LangGraph-inspired stateful workflow engine for Rust.

Topics

Resources

Stars

3 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages