Skip to main content

CompiledGraph

Struct CompiledGraph 

Source
pub struct CompiledGraph { /* private fields */ }
Available on crate feature graph only.
Expand description

A compiled graph ready for execution

Implementations§

Source§

impl CompiledGraph

Convenience methods for CompiledGraph

Source

pub async fn invoke( &self, input: HashMap<String, Value>, config: ExecutionConfig, ) -> Result<HashMap<String, Value>, GraphError>

Execute the graph synchronously

Source

pub async fn invoke_detailed( &self, input: HashMap<String, Value>, config: ExecutionConfig, ) -> Result<GraphOutcome, GraphError>

Executes and reports what the run asked of its caller.

Only a graph run as a SubgraphNode has anything to report beyond its state, so Self::invoke is the usual entry point.

Source

pub fn stream( &self, input: HashMap<String, Value>, config: ExecutionConfig, mode: StreamMode, ) -> impl Stream<Item = Result<StreamEvent, GraphError>>

Execute with streaming

Source

pub async fn get_state( &self, thread_id: &str, ) -> Result<Option<HashMap<String, Value>>, GraphError>

Get current state for a thread

Source

pub async fn update_state( &self, thread_id: &str, updates: impl IntoIterator<Item = (String, Value)>, ) -> Result<(), GraphError>

Update state for a thread (for human-in-the-loop)

Source§

impl CompiledGraph

Source

pub fn with_checkpointer<C>(self, checkpointer: C) -> CompiledGraph
where C: Checkpointer + 'static,

Configure checkpointing

Source

pub fn with_checkpointer_arc( self, checkpointer: Arc<dyn Checkpointer>, ) -> CompiledGraph

Configure checkpointing with Arc

Source

pub fn with_interrupt_before(self, nodes: &[&str]) -> CompiledGraph

Configure interrupt before specific nodes

Source

pub fn with_interrupt_after(self, nodes: &[&str]) -> CompiledGraph

Configure interrupt after specific nodes

Source

pub fn with_recursion_limit(self, limit: usize) -> CompiledGraph

Set recursion limit for cycles

Source

pub fn with_max_concurrency(self, limit: usize) -> CompiledGraph

Cap how many nodes execute concurrently within one super-step.

A wide fan-out otherwise dispatches its whole frontier at once, which can exhaust a connection pool or trip a provider rate limit. Nodes beyond the cap wait for a slot; the dispatch order is the frontier’s, sorted, so it does not depend on timing.

Without this the frontier runs unbounded, which stays the default.

Source

pub fn with_strict_channels(self) -> CompiledGraph

Fail the run when a node writes a channel the schema does not declare.

An undeclared channel otherwise takes the overwrite reducer, because that is the fallback for a name the schema does not hold. A graph that declared a list channel and then wrote a near-miss name keeps only the last value and reports nothing. Enforcement turns that into crate::error::GraphError::UndeclaredChannel.

A graph that declares no channels accepts any name even under enforcement, because there is nothing to check against.

Off by default: a graph may legitimately declare the channels a caller reads and let its nodes pass other values between themselves.

§Example
use adk_graph::edge::{END, START};
use adk_graph::graph::StateGraph;
use adk_graph::node::NodeOutput;
use serde_json::json;

let graph = StateGraph::with_channels(&["total"])
    .add_node_fn("sum", |_ctx| async move {
        Ok(NodeOutput::new().with_update("total", json!(3)))
    })
    .add_edge(START, "sum")
    .add_edge("sum", END)
    .compile()
    .unwrap()
    .with_strict_channels();
Source

pub fn with_checkpoint_retention(self, policy: RetentionPolicy) -> CompiledGraph

Discards old checkpoints as the run proceeds.

A thread otherwise accumulates one checkpoint per super-step for as long as it lives, which costs storage and slows a list. The newest is always kept, because it is the one a resume loads.

Off by default, so an existing thread keeps its whole history and time travel can still reach every step.

§Example
use adk_graph::checkpoint::{MemoryCheckpointer, RetentionPolicy};
use adk_graph::edge::{END, START};
use adk_graph::graph::StateGraph;
use adk_graph::node::NodeOutput;

let graph = StateGraph::with_channels(&["value"])
    .add_node_fn("step", |_ctx| async move { Ok(NodeOutput::new()) })
    .add_edge(START, "step")
    .add_edge("step", END)
    .compile()?
    .with_checkpointer(MemoryCheckpointer::new())
    .with_checkpoint_retention(RetentionPolicy::keep_last(20));
Source

pub fn with_node_defaults(self, defaults: NodeDefaults) -> CompiledGraph

Applies policies to every node that does not set its own.

Repeating the same retry or timeout across twenty nodes is easy to get wrong by omission. A per-node value always wins over the default.

§Example
use adk_graph::edge::{END, START};
use adk_graph::graph::{NodeDefaults, StateGraph};
use adk_graph::node::NodeOutput;
use adk_graph::retry::RetryPolicy;

let graph = StateGraph::with_channels(&["value"])
    .add_node_fn("fetch", |_ctx| async move { Ok(NodeOutput::new()) })
    .add_edge(START, "fetch")
    .add_edge("fetch", END)
    .compile()?
    // Every node retries three times, unless it says otherwise.
    .with_node_defaults(NodeDefaults::new().with_retry(RetryPolicy::new(3)))
    // And this one gets five.
    .with_node_retry("fetch", RetryPolicy::new(5));
Source

pub fn with_node_error_handler<F>(self, node: &str, handler: F) -> CompiledGraph
where F: Fn(&str, &GraphError, &HashMap<String, Value>) -> Result<NodeOutput, GraphError> + Send + Sync + 'static,

Handles one node’s failure instead of ending the run.

Called once the node’s retry budget is spent. The handler receives the node name, the error, and the state as it stands, and returns the updates to apply — typically recording what failed and naming a recovery node with NodeOutput::with_goto. Returning Err ends the run.

An interrupt is never routed here: a pause is not a failure.

Source

pub fn has_checkpointer(&self) -> bool

Whether this graph holds a checkpointer.

Source

pub fn can_pause(&self) -> bool

Whether this graph declares any static interrupt gate.

A dynamic interrupt cannot be seen from the graph, because a node decides at run time, so this reports only the declared gates.

Source

pub fn with_node_retry(self, node: &str, policy: RetryPolicy) -> CompiledGraph

Attach a retry policy to one node.

A node with no policy is attempted once, which is the behaviour of a graph that configures none.

Source

pub fn node(&self, name: &str) -> Option<Arc<dyn Node>>

A node by name, for a caller that needs to run one on its own.

Source

pub fn state_channels(&self) -> Vec<String>

The declared state channel names, sorted.

Source

pub fn timeout_policy_for(&self, node_name: &str) -> Option<&TimeoutPolicy>

Get the effective timeout policy for a node.

Returns the per-node policy if one was configured via GraphAgentBuilder::node_timeout, otherwise falls back to the default timeout policy. Returns None if neither is set.

Source

pub fn get_entry_nodes(&self) -> Vec<String>

Get entry nodes

Source

pub fn get_next_nodes( &self, executed: &[String], state: &HashMap<String, Value>, ) -> Result<Vec<String>, GraphError>

Get next nodes after executing the given nodes

§Errors

Returns GraphError::UnknownRouteTarget when a router answers with a key that is not among the declared targets. A route to END is declared, so it is not an error; a key nobody declared is, because the branch would otherwise stop and the run would report success having skipped the work.

Source

pub fn route_dispatches( &self, executed: &[String], state: &HashMap<String, Value>, ) -> Result<Vec<(String, Vec<String>)>, GraphError>

Reports the conditional dispatches the executed nodes produce.

Only conditional edges appear: a direct edge involves no decision. Used for StreamEvent::RouteDispatched, and called only when a caller asked for the debug stream, so a router is not evaluated again on the common path.

§Errors

Returns GraphError::UnknownRouteTarget on an undeclared route key, matching Self::get_next_nodes.

Source

pub fn leads_to_end( &self, executed: &[String], state: &HashMap<String, Value>, ) -> bool

Check if any of the executed nodes lead to END

Source

pub fn get_upstream_nodes(&self, target_node: &str) -> Vec<String>

Get all upstream source nodes for a given target node.

Returns the names of all nodes that have an edge pointing to the given target node. This is used by the deferred node scheduler to determine which upstream paths must complete before a fan-in node can execute.

For conditional edges, all possible source nodes are included since any of them could route to the target at runtime.

Source

pub fn schema(&self) -> &StateSchema

Get the state schema

Source

pub fn checkpointer(&self) -> Option<&Arc<dyn Checkpointer>>

Get the checkpointer if configured

Source§

impl CompiledGraph

Source

pub fn time_travel( &self, thread_id: &str, ) -> Result<TimeTravelHandle<'_>, GraphError>

Create a TimeTravelHandle for navigating the execution history of a thread.

Requires that the graph was compiled with a checkpointer. If no checkpointer is configured, this method panics.

§Arguments
  • thread_id - The thread identifier whose history to navigate
§Panics

Panics if the graph was compiled without a checkpointer.

§Example
use adk_graph::prelude::*;

let checkpointer = Arc::new(MemoryCheckpointer::new());
let graph = StateGraph::new(StateSchema::simple(&["data"]))
    .add_node("process", process_fn)
    .add_edge(START, "process")
    .add_edge("process", END)
    .compile(Some(checkpointer));

let handle = graph.time_travel("thread_1").unwrap();
let steps = handle.steps().await?;
§Errors

Returns GraphError::CheckpointError when the graph has no checkpointer, because every operation on the handle reads checkpoints.

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<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<T> Conv for T

Source§

fn conv<T>(self) -> T
where Self: Into<T>,

Converts self into T using Into<T>. Read more
Source§

impl<T> FmtForward for T

Source§

fn fmt_binary(self) -> FmtBinary<Self>
where Self: Binary,

Causes self to use its Binary implementation when Debug-formatted.
Source§

fn fmt_display(self) -> FmtDisplay<Self>
where Self: Display,

Causes self to use its Display implementation when Debug-formatted.
Source§

fn fmt_lower_exp(self) -> FmtLowerExp<Self>
where Self: LowerExp,

Causes self to use its LowerExp implementation when Debug-formatted.
Source§

fn fmt_lower_hex(self) -> FmtLowerHex<Self>
where Self: LowerHex,

Causes self to use its LowerHex implementation when Debug-formatted.
Source§

fn fmt_octal(self) -> FmtOctal<Self>
where Self: Octal,

Causes self to use its Octal implementation when Debug-formatted.
Source§

fn fmt_pointer(self) -> FmtPointer<Self>
where Self: Pointer,

Causes self to use its Pointer implementation when Debug-formatted.
Source§

fn fmt_upper_exp(self) -> FmtUpperExp<Self>
where Self: UpperExp,

Causes self to use its UpperExp implementation when Debug-formatted.
Source§

fn fmt_upper_hex(self) -> FmtUpperHex<Self>
where Self: UpperHex,

Causes self to use its UpperHex implementation when Debug-formatted.
Source§

fn fmt_list(self) -> FmtList<Self>
where &'a Self: for<'a> IntoIterator,

Formats each item in a sequence. 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> MaybeSend for T
where T: Send,

Source§

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

Source§

fn pipe<R>(self, func: impl FnOnce(Self) -> R) -> R
where Self: Sized,

Pipes by value. This is generally the method you want to use. Read more
Source§

fn pipe_ref<'a, R>(&'a self, func: impl FnOnce(&'a Self) -> R) -> R
where R: 'a,

Borrows self and passes that borrow into the pipe function. Read more
Source§

fn pipe_ref_mut<'a, R>(&'a mut self, func: impl FnOnce(&'a mut Self) -> R) -> R
where R: 'a,

Mutably borrows self and passes that borrow into the pipe function. Read more
Source§

fn pipe_borrow<'a, B, R>(&'a self, func: impl FnOnce(&'a B) -> R) -> R
where Self: Borrow<B>, B: 'a + ?Sized, R: 'a,

Borrows self, then passes self.borrow() into the pipe function. Read more
Source§

fn pipe_borrow_mut<'a, B, R>( &'a mut self, func: impl FnOnce(&'a mut B) -> R, ) -> R
where Self: BorrowMut<B>, B: 'a + ?Sized, R: 'a,

Mutably borrows self, then passes self.borrow_mut() into the pipe function. Read more
Source§

fn pipe_as_ref<'a, U, R>(&'a self, func: impl FnOnce(&'a U) -> R) -> R
where Self: AsRef<U>, U: 'a + ?Sized, R: 'a,

Borrows self, then passes self.as_ref() into the pipe function.
Source§

fn pipe_as_mut<'a, U, R>(&'a mut self, func: impl FnOnce(&'a mut U) -> R) -> R
where Self: AsMut<U>, U: 'a + ?Sized, R: 'a,

Mutably borrows self, then passes self.as_mut() into the pipe function.
Source§

fn pipe_deref<'a, T, R>(&'a self, func: impl FnOnce(&'a T) -> R) -> R
where Self: Deref<Target = T>, T: 'a + ?Sized, R: 'a,

Borrows self, then passes self.deref() into the pipe function.
Source§

fn pipe_deref_mut<'a, T, R>( &'a mut self, func: impl FnOnce(&'a mut T) -> R, ) -> R
where Self: DerefMut<Target = T> + Deref, T: 'a + ?Sized, R: 'a,

Mutably borrows self, then passes self.deref_mut() into the pipe function.
Source§

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

Source§

fn and<P, B, E>(self, other: P) -> And<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow only if self and other return Action::Follow. Read more
Source§

fn or<P, B, E>(self, other: P) -> Or<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow if either self or other returns Action::Follow. Read more
Source§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T> Tap for T

Source§

fn tap(self, func: impl FnOnce(&Self)) -> Self

Immutable access to a value. Read more
Source§

fn tap_mut(self, func: impl FnOnce(&mut Self)) -> Self

Mutable access to a value. Read more
Source§

fn tap_borrow<B>(self, func: impl FnOnce(&B)) -> Self
where Self: Borrow<B>, B: ?Sized,

Immutable access to the Borrow<B> of a value. Read more
Source§

fn tap_borrow_mut<B>(self, func: impl FnOnce(&mut B)) -> Self
where Self: BorrowMut<B>, B: ?Sized,

Mutable access to the BorrowMut<B> of a value. Read more
Source§

fn tap_ref<R>(self, func: impl FnOnce(&R)) -> Self
where Self: AsRef<R>, R: ?Sized,

Immutable access to the AsRef<R> view of a value. Read more
Source§

fn tap_ref_mut<R>(self, func: impl FnOnce(&mut R)) -> Self
where Self: AsMut<R>, R: ?Sized,

Mutable access to the AsMut<R> view of a value. Read more
Source§

fn tap_deref<T>(self, func: impl FnOnce(&T)) -> Self
where Self: Deref<Target = T>, T: ?Sized,

Immutable access to the Deref::Target of a value. Read more
Source§

fn tap_deref_mut<T>(self, func: impl FnOnce(&mut T)) -> Self
where Self: DerefMut<Target = T> + Deref, T: ?Sized,

Mutable access to the Deref::Target of a value. Read more
Source§

fn tap_dbg(self, func: impl FnOnce(&Self)) -> Self

Calls .tap() only in debug builds, and is erased in release builds.
Source§

fn tap_mut_dbg(self, func: impl FnOnce(&mut Self)) -> Self

Calls .tap_mut() only in debug builds, and is erased in release builds.
Source§

fn tap_borrow_dbg<B>(self, func: impl FnOnce(&B)) -> Self
where Self: Borrow<B>, B: ?Sized,

Calls .tap_borrow() only in debug builds, and is erased in release builds.
Source§

fn tap_borrow_mut_dbg<B>(self, func: impl FnOnce(&mut B)) -> Self
where Self: BorrowMut<B>, B: ?Sized,

Calls .tap_borrow_mut() only in debug builds, and is erased in release builds.
Source§

fn tap_ref_dbg<R>(self, func: impl FnOnce(&R)) -> Self
where Self: AsRef<R>, R: ?Sized,

Calls .tap_ref() only in debug builds, and is erased in release builds.
Source§

fn tap_ref_mut_dbg<R>(self, func: impl FnOnce(&mut R)) -> Self
where Self: AsMut<R>, R: ?Sized,

Calls .tap_ref_mut() only in debug builds, and is erased in release builds.
Source§

fn tap_deref_dbg<T>(self, func: impl FnOnce(&T)) -> Self
where Self: Deref<Target = T>, T: ?Sized,

Calls .tap_deref() only in debug builds, and is erased in release builds.
Source§

fn tap_deref_mut_dbg<T>(self, func: impl FnOnce(&mut T)) -> Self
where Self: DerefMut<Target = T> + Deref, T: ?Sized,

Calls .tap_deref_mut() only in debug builds, and is erased in release builds.
Source§

impl<T> TryConv for T

Source§

fn try_conv<T>(self) -> Result<T, Self::Error>
where Self: TryInto<T>,

Attempts to convert self into T using TryInto<T>. 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