pub struct CompiledGraph { /* private fields */ }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: State,
config: ExecutionConfig,
) -> Result<State>
pub async fn invoke( &self, input: State, config: ExecutionConfig, ) -> Result<State>
Execute the graph synchronously
Sourcepub async fn invoke_detailed(
&self,
input: State,
config: ExecutionConfig,
) -> Result<GraphOutcome>
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.
Sourcepub fn stream(
&self,
input: State,
config: ExecutionConfig,
mode: StreamMode,
) -> impl Stream<Item = Result<StreamEvent>> + '_
pub fn stream( &self, input: State, config: ExecutionConfig, mode: StreamMode, ) -> impl Stream<Item = Result<StreamEvent>> + '_
Execute with streaming
Sourcepub async fn get_state(&self, thread_id: &str) -> Result<Option<State>>
pub async fn get_state(&self, thread_id: &str) -> Result<Option<State>>
Get current state for a thread
Sourcepub async fn update_state(
&self,
thread_id: &str,
updates: impl IntoIterator<Item = (String, Value)>,
) -> Result<()>
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
impl CompiledGraph
Sourcepub fn with_checkpointer<C: Checkpointer + 'static>(
self,
checkpointer: C,
) -> Self
pub fn with_checkpointer<C: Checkpointer + 'static>( self, checkpointer: C, ) -> Self
Configure checkpointing
Sourcepub fn with_checkpointer_arc(self, checkpointer: Arc<dyn Checkpointer>) -> Self
pub fn with_checkpointer_arc(self, checkpointer: Arc<dyn Checkpointer>) -> Self
Configure checkpointing with Arc
Sourcepub fn with_interrupt_before(self, nodes: &[&str]) -> Self
pub fn with_interrupt_before(self, nodes: &[&str]) -> Self
Configure interrupt before specific nodes
Sourcepub fn with_interrupt_after(self, nodes: &[&str]) -> Self
pub fn with_interrupt_after(self, nodes: &[&str]) -> Self
Configure interrupt after specific nodes
Sourcepub fn with_recursion_limit(self, limit: usize) -> Self
pub fn with_recursion_limit(self, limit: usize) -> Self
Set recursion limit for cycles
Sourcepub fn with_max_concurrency(self, limit: usize) -> Self
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.
Sourcepub fn with_strict_channels(self) -> Self
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();Sourcepub fn with_checkpoint_retention(self, policy: RetentionPolicy) -> Self
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));Sourcepub fn with_node_defaults(self, defaults: NodeDefaults) -> Self
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));Sourcepub fn with_node_error_handler<F>(self, node: &str, handler: F) -> Self
pub fn with_node_error_handler<F>(self, node: &str, handler: F) -> Self
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) -> Self
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.
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: &State,
) -> Result<Vec<String>>
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.
Sourcepub fn route_dispatches(
&self,
executed: &[String],
state: &State,
) -> Result<Vec<(String, Vec<String>)>>
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.
Sourcepub fn leads_to_end(&self, executed: &[String], state: &State) -> bool
pub fn leads_to_end(&self, executed: &[String], state: &State) -> 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