pub struct CompiledGraph { /* private fields */ }graph only.Expand description
A compiled graph ready for execution
Implementations§
Source§impl CompiledGraph
Convenience methods for CompiledGraph
impl CompiledGraph
Convenience methods for CompiledGraph
Sourcepub async fn invoke(
&self,
input: HashMap<String, Value>,
config: ExecutionConfig,
) -> Result<HashMap<String, Value>, GraphError>
pub async fn invoke( &self, input: HashMap<String, Value>, config: ExecutionConfig, ) -> Result<HashMap<String, Value>, GraphError>
Execute the graph synchronously
Sourcepub async fn invoke_detailed(
&self,
input: HashMap<String, Value>,
config: ExecutionConfig,
) -> Result<GraphOutcome, GraphError>
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.
Sourcepub fn stream(
&self,
input: HashMap<String, Value>,
config: ExecutionConfig,
mode: StreamMode,
) -> impl Stream<Item = Result<StreamEvent, GraphError>>
pub fn stream( &self, input: HashMap<String, Value>, config: ExecutionConfig, mode: StreamMode, ) -> impl Stream<Item = Result<StreamEvent, GraphError>>
Execute with streaming
Sourcepub async fn get_state(
&self,
thread_id: &str,
) -> Result<Option<HashMap<String, Value>>, GraphError>
pub async fn get_state( &self, thread_id: &str, ) -> Result<Option<HashMap<String, Value>>, GraphError>
Get current state for a thread
Sourcepub async fn update_state(
&self,
thread_id: &str,
updates: impl IntoIterator<Item = (String, Value)>,
) -> Result<(), GraphError>
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
impl CompiledGraph
Sourcepub fn with_checkpointer<C>(self, checkpointer: C) -> CompiledGraphwhere
C: Checkpointer + 'static,
pub fn with_checkpointer<C>(self, checkpointer: C) -> CompiledGraphwhere
C: Checkpointer + 'static,
Configure checkpointing
Sourcepub fn with_checkpointer_arc(
self,
checkpointer: Arc<dyn Checkpointer>,
) -> CompiledGraph
pub fn with_checkpointer_arc( self, checkpointer: Arc<dyn Checkpointer>, ) -> CompiledGraph
Configure checkpointing with Arc
Sourcepub fn with_interrupt_before(self, nodes: &[&str]) -> CompiledGraph
pub fn with_interrupt_before(self, nodes: &[&str]) -> CompiledGraph
Configure interrupt before specific nodes
Sourcepub fn with_interrupt_after(self, nodes: &[&str]) -> CompiledGraph
pub fn with_interrupt_after(self, nodes: &[&str]) -> CompiledGraph
Configure interrupt after specific nodes
Sourcepub fn with_recursion_limit(self, limit: usize) -> CompiledGraph
pub fn with_recursion_limit(self, limit: usize) -> CompiledGraph
Set recursion limit for cycles
Sourcepub fn with_max_concurrency(self, limit: usize) -> CompiledGraph
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.
Sourcepub fn with_strict_channels(self) -> CompiledGraph
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();Sourcepub fn with_checkpoint_retention(self, policy: RetentionPolicy) -> CompiledGraph
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));Sourcepub fn with_node_defaults(self, defaults: NodeDefaults) -> CompiledGraph
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));Sourcepub fn with_node_error_handler<F>(self, node: &str, handler: F) -> CompiledGraphwhere
F: Fn(&str, &GraphError, &HashMap<String, Value>) -> Result<NodeOutput, GraphError> + Send + Sync + 'static,
pub fn with_node_error_handler<F>(self, node: &str, handler: F) -> CompiledGraphwhere
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.
Sourcepub fn has_checkpointer(&self) -> bool
pub fn has_checkpointer(&self) -> bool
Whether this graph holds a checkpointer.
Sourcepub fn can_pause(&self) -> bool
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.
Sourcepub fn with_node_retry(self, node: &str, policy: RetryPolicy) -> CompiledGraph
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.
Sourcepub fn node(&self, name: &str) -> Option<Arc<dyn Node>>
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.
Sourcepub fn state_channels(&self) -> Vec<String>
pub fn state_channels(&self) -> Vec<String>
The declared state channel names, sorted.
Sourcepub fn timeout_policy_for(&self, node_name: &str) -> Option<&TimeoutPolicy>
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.
Sourcepub fn get_entry_nodes(&self) -> Vec<String>
pub fn get_entry_nodes(&self) -> Vec<String>
Get entry nodes
Sourcepub fn get_next_nodes(
&self,
executed: &[String],
state: &HashMap<String, Value>,
) -> Result<Vec<String>, GraphError>
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.
Sourcepub fn route_dispatches(
&self,
executed: &[String],
state: &HashMap<String, Value>,
) -> Result<Vec<(String, Vec<String>)>, GraphError>
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.
Sourcepub fn leads_to_end(
&self,
executed: &[String],
state: &HashMap<String, Value>,
) -> bool
pub fn leads_to_end( &self, executed: &[String], state: &HashMap<String, Value>, ) -> bool
Check if any of the executed nodes lead to END
Sourcepub fn get_upstream_nodes(&self, target_node: &str) -> Vec<String>
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.
Sourcepub fn schema(&self) -> &StateSchema
pub fn schema(&self) -> &StateSchema
Get the state schema
Sourcepub fn checkpointer(&self) -> Option<&Arc<dyn Checkpointer>>
pub fn checkpointer(&self) -> Option<&Arc<dyn Checkpointer>>
Get the checkpointer if configured
Source§impl CompiledGraph
impl CompiledGraph
Sourcepub fn time_travel(
&self,
thread_id: &str,
) -> Result<TimeTravelHandle<'_>, GraphError>
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§
impl !RefUnwindSafe for CompiledGraph
impl !UnwindSafe for CompiledGraph
impl Freeze for CompiledGraph
impl Send for CompiledGraph
impl Sync for CompiledGraph
impl Unpin for CompiledGraph
impl UnsafeUnpin for CompiledGraph
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
Source§impl<T> FmtForward for T
impl<T> FmtForward for T
Source§fn fmt_binary(self) -> FmtBinary<Self>where
Self: Binary,
fn fmt_binary(self) -> FmtBinary<Self>where
Self: Binary,
self to use its Binary implementation when Debug-formatted.Source§fn fmt_display(self) -> FmtDisplay<Self>where
Self: Display,
fn fmt_display(self) -> FmtDisplay<Self>where
Self: Display,
self to use its Display implementation when
Debug-formatted.Source§fn fmt_lower_exp(self) -> FmtLowerExp<Self>where
Self: LowerExp,
fn fmt_lower_exp(self) -> FmtLowerExp<Self>where
Self: LowerExp,
self to use its LowerExp implementation when
Debug-formatted.Source§fn fmt_lower_hex(self) -> FmtLowerHex<Self>where
Self: LowerHex,
fn fmt_lower_hex(self) -> FmtLowerHex<Self>where
Self: LowerHex,
self to use its LowerHex implementation when
Debug-formatted.Source§fn fmt_octal(self) -> FmtOctal<Self>where
Self: Octal,
fn fmt_octal(self) -> FmtOctal<Self>where
Self: Octal,
self to use its Octal implementation when Debug-formatted.Source§fn fmt_pointer(self) -> FmtPointer<Self>where
Self: Pointer,
fn fmt_pointer(self) -> FmtPointer<Self>where
Self: Pointer,
self to use its Pointer implementation when
Debug-formatted.Source§fn fmt_upper_exp(self) -> FmtUpperExp<Self>where
Self: UpperExp,
fn fmt_upper_exp(self) -> FmtUpperExp<Self>where
Self: UpperExp,
self to use its UpperExp implementation when
Debug-formatted.Source§fn fmt_upper_hex(self) -> FmtUpperHex<Self>where
Self: UpperHex,
fn fmt_upper_hex(self) -> FmtUpperHex<Self>where
Self: UpperHex,
self to use its UpperHex implementation when
Debug-formatted.Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
Source§fn in_current_span(self) -> Instrumented<Self> ⓘ
fn in_current_span(self) -> Instrumented<Self> ⓘ
Source§impl<T> IntoEither for T
impl<T> IntoEither for T
Source§fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
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 moreSource§fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
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 moreimpl<T> MaybeSend for Twhere
T: Send,
Source§impl<T> Pipe for Twhere
T: ?Sized,
impl<T> Pipe for Twhere
T: ?Sized,
Source§fn pipe<R>(self, func: impl FnOnce(Self) -> R) -> Rwhere
Self: Sized,
fn pipe<R>(self, func: impl FnOnce(Self) -> R) -> Rwhere
Self: Sized,
Source§fn pipe_ref<'a, R>(&'a self, func: impl FnOnce(&'a Self) -> R) -> Rwhere
R: 'a,
fn pipe_ref<'a, R>(&'a self, func: impl FnOnce(&'a Self) -> R) -> Rwhere
R: 'a,
self and passes that borrow into the pipe function. Read moreSource§fn pipe_ref_mut<'a, R>(&'a mut self, func: impl FnOnce(&'a mut Self) -> R) -> Rwhere
R: 'a,
fn pipe_ref_mut<'a, R>(&'a mut self, func: impl FnOnce(&'a mut Self) -> R) -> Rwhere
R: 'a,
self and passes that borrow into the pipe function. Read moreSource§fn pipe_borrow<'a, B, R>(&'a self, func: impl FnOnce(&'a B) -> R) -> R
fn pipe_borrow<'a, B, R>(&'a self, func: impl FnOnce(&'a B) -> R) -> R
Source§fn pipe_borrow_mut<'a, B, R>(
&'a mut self,
func: impl FnOnce(&'a mut B) -> R,
) -> R
fn pipe_borrow_mut<'a, B, R>( &'a mut self, func: impl FnOnce(&'a mut B) -> R, ) -> R
Source§fn pipe_as_ref<'a, U, R>(&'a self, func: impl FnOnce(&'a U) -> R) -> R
fn pipe_as_ref<'a, U, R>(&'a self, func: impl FnOnce(&'a U) -> R) -> R
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
fn pipe_as_mut<'a, U, R>(&'a mut self, func: impl FnOnce(&'a mut U) -> R) -> R
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
fn pipe_deref<'a, T, R>(&'a self, func: impl FnOnce(&'a T) -> R) -> R
self, then passes self.deref() into the pipe function.Source§impl<T> PolicyExt for Twhere
T: ?Sized,
impl<T> PolicyExt for Twhere
T: ?Sized,
impl<T> Read<Exclusive, BecauseExclusive> for Twhere
T: ?Sized,
Source§impl<T> Tap for T
impl<T> Tap for T
Source§fn tap_borrow<B>(self, func: impl FnOnce(&B)) -> Self
fn tap_borrow<B>(self, func: impl FnOnce(&B)) -> Self
Borrow<B> of a value. Read moreSource§fn tap_borrow_mut<B>(self, func: impl FnOnce(&mut B)) -> Self
fn tap_borrow_mut<B>(self, func: impl FnOnce(&mut B)) -> Self
BorrowMut<B> of a value. Read moreSource§fn tap_ref<R>(self, func: impl FnOnce(&R)) -> Self
fn tap_ref<R>(self, func: impl FnOnce(&R)) -> Self
AsRef<R> view of a value. Read moreSource§fn tap_ref_mut<R>(self, func: impl FnOnce(&mut R)) -> Self
fn tap_ref_mut<R>(self, func: impl FnOnce(&mut R)) -> Self
AsMut<R> view of a value. Read moreSource§fn tap_deref<T>(self, func: impl FnOnce(&T)) -> Self
fn tap_deref<T>(self, func: impl FnOnce(&T)) -> Self
Deref::Target of a value. Read moreSource§fn tap_deref_mut<T>(self, func: impl FnOnce(&mut T)) -> Self
fn tap_deref_mut<T>(self, func: impl FnOnce(&mut T)) -> Self
Deref::Target of a value. Read moreSource§fn tap_dbg(self, func: impl FnOnce(&Self)) -> Self
fn tap_dbg(self, func: impl FnOnce(&Self)) -> Self
.tap() only in debug builds, and is erased in release builds.Source§fn tap_mut_dbg(self, func: impl FnOnce(&mut Self)) -> Self
fn tap_mut_dbg(self, func: impl FnOnce(&mut Self)) -> Self
.tap_mut() only in debug builds, and is erased in release
builds.Source§fn tap_borrow_dbg<B>(self, func: impl FnOnce(&B)) -> Self
fn tap_borrow_dbg<B>(self, func: impl FnOnce(&B)) -> Self
.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
fn tap_borrow_mut_dbg<B>(self, func: impl FnOnce(&mut B)) -> Self
.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
fn tap_ref_dbg<R>(self, func: impl FnOnce(&R)) -> Self
.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
fn tap_ref_mut_dbg<R>(self, func: impl FnOnce(&mut R)) -> Self
.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
fn tap_deref_dbg<T>(self, func: impl FnOnce(&T)) -> Self
.tap_deref() only in debug builds, and is erased in release
builds.