Skip to main content

FlowRunner

Struct FlowRunner 

Source
pub struct FlowRunner { /* private fields */ }
Expand description

High-level helper that orchestrates the common load → execute → save pattern.

FlowRunner provides a convenient wrapper around the lower-level graph execution API. It automatically handles session loading, execution, and persistence.

§When to Use FlowRunner

  • Web services: Execute one step per HTTP request
  • Interactive applications: Step-by-step workflow progression
  • Simple demos: Minimal boilerplate for common use cases

§Performance

FlowRunner is lightweight and efficient:

  • Creation cost: ~2 pointer copies (negligible)
  • Memory overhead: 16 bytes (2 × Arc<T>)
  • Runtime cost: Identical to manual approach

§Examples

§Basic Usage

use graph_flow::{FlowRunner, Graph, InMemorySessionStorage, Session, SessionStorage};
use std::sync::Arc;

let graph = Arc::new(Graph::new("my_workflow"));
let storage = Arc::new(InMemorySessionStorage::new());
let runner = FlowRunner::new(graph, storage.clone());

// Create a session first
let session = Session::new_from_task("session_id".to_string(), "start_task");
storage.save(session).await?;

// Execute workflow step
let result = runner.run("session_id").await?;
println!("Response: {:?}", result.response);
use graph_flow::FlowRunner;
use std::sync::Arc;

// Application state
struct AppState {
    flow_runner: Arc<FlowRunner>,
}

impl AppState {
    fn new(runner: FlowRunner) -> Self {
        Self {
            flow_runner: Arc::new(runner),
        }
    }
}

// Request handler
async fn handle_request(
    state: Arc<AppState>,
    session_id: String,
) -> Result<String, Box<dyn std::error::Error>> {
    let result = state.flow_runner.run(&session_id).await?;
    Ok(result.response.unwrap_or_default())
}

Implementations§

Source§

impl FlowRunner

Source

pub fn new(graph: Arc<Graph>, storage: Arc<dyn SessionStorage>) -> Self

Create a new FlowRunner from an Arc<Graph> and any SessionStorage implementation.

§Parameters
  • graph - The workflow graph to execute
  • storage - Storage backend for session persistence
§Examples
use graph_flow::{FlowRunner, Graph, InMemorySessionStorage};
use std::sync::Arc;

let graph = Arc::new(Graph::new("my_workflow"));
let storage = Arc::new(InMemorySessionStorage::new());
let runner = FlowRunner::new(graph, storage);
§With PostgreSQL Storage
use graph_flow::{FlowRunner, Graph, PostgresSessionStorage};
use std::sync::Arc;

let graph = Arc::new(Graph::new("my_workflow"));
let storage = Arc::new(
    PostgresSessionStorage::connect("postgresql://localhost/mydb").await?
);
let runner = FlowRunner::new(graph, storage);
Source

pub async fn run(&self, session_id: &str) -> Result<ExecutionResult>

Execute exactly one task for the given session_id and persist the updated session.

This method:

  1. Loads the session from storage
  2. Executes the current task
  3. Saves the updated session back to storage
  4. Returns the execution result
§Parameters
  • session_id - Unique identifier for the session to execute
§Returns

Returns the same ExecutionResult that Graph::execute_session does, so callers can inspect the assistant’s response and the status (WaitingForInput, Completed, etc.).

§Errors

Returns an error if:

  • The session doesn’t exist
  • Task execution fails
  • Storage operations fail
§Concurrency

Sessions are protected by optimistic locking: if another run call (or any other writer) saved the same session between this call’s load and save, the save fails with crate::GraphError::SessionConflict instead of silently losing the other writer’s update. On conflict, retry by calling run again (it reloads the latest session state). Note that the task’s side effects from the conflicting attempt are not rolled back, so tasks should be idempotent if concurrent runs are possible.

§Examples
§Basic Execution
let result = runner.run("test_session").await?;

match result.status {
    graph_flow::ExecutionStatus::Completed => {
        println!("Workflow completed: {:?}", result.response);
    }
    graph_flow::ExecutionStatus::WaitingForInput => {
        println!("Waiting for user input: {:?}", result.response);
    }
    graph_flow::ExecutionStatus::Paused { next_task_id, reason } => {
        println!("Paused, next task: {}, reason: {}", next_task_id, reason);
    }
}
§Interactive Loop
loop {
    let result = runner.run("session_id").await?;
     
    match result.status {
        ExecutionStatus::Completed => break,
        ExecutionStatus::WaitingForInput => {
            // Get user input and update context
            // Then continue loop
            break; // For demo
        }
        ExecutionStatus::Paused { .. } => {
            // Continue to next step
            continue;
        }
    }
}
§Error Handling
match runner.run("nonexistent_session").await {
    Ok(result) => {
        println!("Success: {:?}", result.response);
    }
    Err(GraphError::SessionNotFound(session_id)) => {
        eprintln!("Session not found: {}", session_id);
    }
    Err(GraphError::TaskExecutionFailed(msg)) => {
        eprintln!("Task failed: {}", msg);
    }
    Err(e) => {
        eprintln!("Other error: {}", e);
    }
}

Trait Implementations§

Source§

impl Clone for FlowRunner

Source§

fn clone(&self) -> FlowRunner

Returns a duplicate of the value. Read more
1.0.0 (const: unstable) · Source§

fn clone_from(&mut self, source: &Self)

Performs copy-assignment from source. Read more

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> CloneToUninit for T
where T: Clone,

Source§

unsafe fn clone_to_uninit(&self, dest: *mut u8)

🔬This is a nightly-only experimental API. (clone_to_uninit)
Performs copy-assignment from self to dest. Read more
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self>

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self>

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> IntoEither for T

Source§

fn into_either(self, into_left: bool) -> Either<Self, Self>

Converts self into a Left variant of Either<Self, Self> if into_left is true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
where F: FnOnce(&Self) -> bool,

Converts self into a Left variant of Either<Self, Self> if into_left(&self) returns true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T> ToOwned for T
where T: Clone,

Source§

type Owned = T

The resulting type after obtaining ownership.
Source§

fn to_owned(&self) -> T

Creates owned data from borrowed data, usually by cloning. Read more
Source§

fn clone_into(&self, target: &mut T)

Uses borrowed data to replace owned data, usually by cloning. Read more
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = Infallible

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

Source§

fn vzip(self) -> V

Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self>
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self>

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more