Skip to main content

CompiledGraph

Struct CompiledGraph 

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

A compiled graph ready for execution

Implementations§

Source§

impl CompiledGraph

Convenience methods for CompiledGraph

Source

pub async fn invoke( &self, input: State, config: ExecutionConfig, ) -> Result<State>

Execute the graph synchronously

Source

pub async fn invoke_detailed( &self, input: State, config: ExecutionConfig, ) -> Result<GraphOutcome>

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: State, config: ExecutionConfig, mode: StreamMode, ) -> impl Stream<Item = Result<StreamEvent>> + '_

Execute with streaming

Source

pub async fn get_state(&self, thread_id: &str) -> Result<Option<State>>

Get current state for a thread

Source

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

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

Source§

impl CompiledGraph

Source

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

Configure checkpointing

Source

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

Configure checkpointing with Arc

Source

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

Configure interrupt before specific nodes

Source

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

Configure interrupt after specific nodes

Source

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

Set recursion limit for cycles

Source

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

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

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

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

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) -> Self
where F: Fn(&str, &GraphError, &State) -> Result<NodeOutput> + 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) -> Self

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: &State, ) -> Result<Vec<String>>

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: &State, ) -> Result<Vec<(String, Vec<String>)>>

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: &State) -> 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

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> 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, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

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