Agentes de Grafo
Crie fluxos de trabalho complexos e com estado usando orquestraΓ§Γ£o no estilo LangGraph com integraΓ§Γ£o nativa ao ADK-Rust.
VisΓ£o geral
GraphAgent permite definir fluxos de trabalho como grafos direcionados com nΓ³s e arestas, oferecendo suporte a:
- AgentNode: Encapsule agentes LLM como nΓ³s do grafo com mapeadores personalizados de entrada/saΓda
- Fluxos de Trabalho CΓclicos: Suporte nativo a loops e raciocΓnio iterativo (padrΓ£o ReAct)
- Roteamento Condicional: Roteamento dinΓ’mico de arestas com base no estado
- Gerenciamento de Estado: Estado tipado com redutores (substituiΓ§Γ£o, acrΓ©scimo, soma, personalizado)
- CriaΓ§Γ£o de Pontos de VerificaΓ§Γ£o: Estado persistente para tolerΓ’ncia a falhas e interaΓ§Γ£o humana no processo
- Streaming: VΓ‘rios modos de fluxo (valores, atualizaΓ§Γ΅es, mensagens, depuraΓ§Γ£o)
O crate adk-graph fornece orquestraΓ§Γ£o de fluxos de trabalho no estilo LangGraph para criar fluxos de trabalho de agentes complexos e com estado. Ele traz recursos de fluxos de trabalho baseados em grafos para o ecossistema ADK-Rust, mantendo compatibilidade total com o sistema de agentes do ADK.
Principais benefΓcios:
- Design Visual de Fluxos de Trabalho: Defina lΓ³gicas complexas como grafos intuitivos de nΓ³s e arestas
- ExecuΓ§Γ£o Paralela: VΓ‘rios nΓ³s podem ser executados simultaneamente para obter melhor desempenho
- PersistΓͺncia de Estado: CriaΓ§Γ£o integrada de pontos de verificaΓ§Γ£o para tolerΓ’ncia a falhas e interaΓ§Γ£o humana no processo
- IntegraΓ§Γ£o com LLM: Suporte nativo para encapsular agentes ADK como nΓ³s do grafo
- Roteamento FlexΓvel: Arestas estΓ‘ticas, roteamento condicional e tomada de decisΓ΅es dinΓ’mica
Escolhendo entre os Agentes de Fluxo de Trabalho e o Grafo
ADK-Rust oferece duas formas de orquestrar. Nenhuma substitui a outra, e ambas sΓ£o mantidas.
| VocΓͺ precisa | Use | Por quΓͺ |
|---|---|---|
| Uma ordem fixa de etapas | SequentialAgent | A topologia Γ© a lista. Nada a declarar. |
| VΓ‘rios agentes na mesma entrada | ParallelAgent | DistribuiΓ§Γ£o sem junΓ§Γ£o para configurar. |
| Repetir atΓ© que uma condiΓ§Γ£o seja atendida | LoopAgent | A condiΓ§Γ£o de saΓda Γ© um callback, nΓ£o uma aresta. |
| Um ramo escolhido em tempo de execuΓ§Γ£o | Grafo | Arestas condicionais ou um nΓ³ que nomeia seu prΓ³prio sucessor. |
| Ciclos com um orΓ§amento de etapas | Grafo | recursion_limit limita as superetapas. |
| Uma pausa Γ qual uma pessoa responde mais tarde | Grafo | As interrupΓ§Γ΅es criam um checkpoint para a execuΓ§Γ£o e a retomam. |
| SobrevivΓͺncia apΓ³s a reinicializaΓ§Γ£o de um processo | Grafo | SqliteCheckpointer persiste a cada superetapa. |
| Retrocesso para uma etapa anterior | Grafo | A viagem no tempo cria uma bifurcaΓ§Γ£o de um ponto de verificaΓ§Γ£o. |
Prefira os agentes de fluxo de trabalho quando a estrutura do trabalho for conhecida e linear: eles sΓ£o mais curtos de escrever e nΓ£o hΓ‘ um esquema de estado para manter. Opte pelo grafo quando o fluxo de controle depender dos resultados ou quando uma execuΓ§Γ£o precisar sobreviver ao processo.
Eles sΓ£o combinΓ‘veis
Os trΓͺs agentes de fluxo de trabalho implementam Agent, e AgentNode encapsula qualquer Agent, portanto, um agente de fluxo de trabalho Γ© um nΓ³ do grafo:
use adk_agent::SequentialAgent;
use adk_graph::node::AgentNode;
use std::sync::Arc;
let pipeline = Arc::new(SequentialAgent::new("pipeline", vec![extract, validate]));
let node = AgentNode::new(pipeline as Arc<dyn adk_core::Agent>);
// `node` now goes into a StateGraph like any other node.
GraphAgent tambΓ©m implementa Agent, portanto, o inverso tambΓ©m Γ© vΓ‘lido: um grafo pode ser um subagente de um SequentialAgent. Use o grafo para a parte que precisa de ramificaΓ§Γ£o ou durabilidade, e os agentes de fluxo de trabalho para as partes que nΓ£o precisam.
O que vocΓͺ vai criar
Neste guia, vocΓͺ criarΓ‘ um Pipeline de Processamento de Texto que executa traduΓ§Γ£o e sumarizaΓ§Γ£o em paralelo:
βββββββββββββββββββββββ
User Input β β
βββββββββββββββββΆ β START β
β β
ββββββββββββ¬βββββββββββ
β
βββββββββββββββββ΄ββββββββββββββββ
β β
βΌ βΌ
ββββββββββββββββββββ ββββββββββββββββββββ
β TRANSLATOR β β SUMMARIZER β
β β β β
β π«π· French β β π One sentence β
β Translation β β Summary β
βββββββββββ¬βββββββββ βββββββββββ¬βββββββββ
β β
βββββββββββββββββ¬ββββββββββββββββ
β
βΌ
βββββββββββββββββββββββ
β COMBINE β
β β
β π Merge Results β
ββββββββββββ¬βββββββββββ
β
βΌ
βββββββββββββββββββββββ
β END β
β β
β β
Complete β
βββββββββββββββββββββββ
Conceitos principais:
- NΓ³s - Unidades de processamento que realizam tarefas (agentes LLM, funΓ§Γ΅es ou lΓ³gica personalizada)
- Arestas - Fluxo de controle entre nΓ³s (conexΓ΅es estΓ‘ticas ou roteamento condicional)
- Estado - Dados compartilhados que fluem pelo grafo e persistem entre os nΓ³s
- ExecuΓ§Γ£o paralela - VΓ‘rios nΓ³s podem ser executados simultaneamente para melhorar o desempenho
Entendendo os componentes principais
π§ NΓ³s: os trabalhadores Γ nos nΓ³s que o trabalho efetivamente acontece. Cada nΓ³ pode:
- AgentNode: encapsular um agente LLM para processar linguagem natural
- NΓ³ de funΓ§Γ£o: executar cΓ³digo Rust personalizado para processamento de dados
- NΓ³s integrados: usar lΓ³gica predefinida, como contadores ou validadores
Pense nos nΓ³s como trabalhadores especializados em uma linha de montagem β cada um tem uma tarefa e uma especialidade especΓficas.
π Arestas: o controle do fluxo As arestas determinam como a execuΓ§Γ£o se move pelo seu grafo:
- Arestas estΓ‘ticas: conexΓ΅es diretas (
A β B β C) - Arestas condicionais: roteamento dinΓ’mico com base no estado (
if sentiment == "positive" β positive_handler) - Arestas paralelas: vΓ‘rios caminhos a partir de um nΓ³ (
START β [translator, summarizer])
As arestas sΓ£o como sinais de trΓ’nsito e placas rodoviΓ‘rias que direcionam o fluxo de trabalho.
πΎ Estado: a memΓ³ria compartilhada O estado Γ© um armazenamento de chave-valor do qual todos os nΓ³s podem ler e no qual podem gravar:
- Dados de entrada: informaΓ§Γ΅es iniciais fornecidas ao grafo
- Resultados intermediΓ‘rios: a saΓda de um nΓ³ se torna a entrada de outro
- SaΓda final: o resultado concluΓdo apΓ³s todo o processamento
O estado funciona como um quadro branco compartilhado, onde os nΓ³s podem deixar informaΓ§Γ΅es para que outros as utilizem.
β‘ ExecuΓ§Γ£o paralela: o aumento de velocidade Quando vΓ‘rias arestas saem de um nΓ³, os nΓ³s de destino sΓ£o executados simultaneamente:
- Processamento mais rΓ‘pido: tarefas independentes sΓ£o executadas ao mesmo tempo
- EficiΓͺncia de recursos: melhor utilizaΓ§Γ£o de CPU e E/S
- Escalabilidade: lidar com fluxos de trabalho mais complexos sem uma desaceleraΓ§Γ£o linear
Γ como ter vΓ‘rios trabalhadores executando simultaneamente diferentes partes de um trabalho, em vez de esperar na fila.
InΓcio rΓ‘pido
1. Crie seu projeto
cargo new graph_demo
cd graph_demo
Adicione dependΓͺncias a Cargo.toml:
[dependencies]
adk-graph = { version = "2.1.0", features = ["sqlite"] }
adk-agent = "2.1.0"
adk-model = "2.1.0"
adk-core = "2.1.0"
tokio = { version = "1", features = ["full"] }
dotenvy = "0.15"
serde_json = "1.0"
Crie .env com sua chave API:
echo 'GOOGLE_API_KEY=your-api-key' > .env
2. Exemplo de processamento paralelo
Este Γ© um exemplo completo e funcional que processa texto em paralelo:
use adk_agent::LlmAgentBuilder;
use adk_graph::{
agent::GraphAgent,
edge::{END, START},
node::{AgentNode, ExecutionConfig, NodeOutput},
state::State,
};
use adk_model::GeminiModel;
use serde_json::json;
use std::sync::Arc;
#[tokio::main]
async fn main() -> anyhow::Result<()> {
dotenvy::dotenv().ok();
let api_key = std::env::var("GOOGLE_API_KEY")?;
let model = Arc::new(GeminiModel::new(&api_key, "gemini-3.7-flash")?);
// Create specialized LLM agents
let translator_agent = Arc::new(
LlmAgentBuilder::new("translator")
.description("Translates text to French")
.model(model.clone())
.instruction("Translate the input text to French. Only output the translation.")
.build()?,
);
let summarizer_agent = Arc::new(
LlmAgentBuilder::new("summarizer")
.description("Summarizes text")
.model(model.clone())
.instruction("Summarize the input text in one sentence.")
.build()?,
);
// Wrap agents as graph nodes with input/output mappers
let translator_node = AgentNode::new(translator_agent)
.with_input_mapper(|state| {
let text = state.get("input").and_then(|v| v.as_str()).unwrap_or("");
adk_core::Content::new("user").with_text(text)
})
.with_output_mapper(|events| {
let mut updates = std::collections::HashMap::new();
for event in events {
if let Some(content) = event.content() {
let text: String = content.parts.iter()
.filter_map(|p| p.text())
.collect::<Vec<_>>()
.join("");
if !text.is_empty() {
updates.insert("translation".to_string(), json!(text));
}
}
}
updates
});
let summarizer_node = AgentNode::new(summarizer_agent)
.with_input_mapper(|state| {
let text = state.get("input").and_then(|v| v.as_str()).unwrap_or("");
adk_core::Content::new("user").with_text(text)
})
.with_output_mapper(|events| {
let mut updates = std::collections::HashMap::new();
for event in events {
if let Some(content) = event.content() {
let text: String = content.parts.iter()
.filter_map(|p| p.text())
.collect::<Vec<_>>()
.join("");
if !text.is_empty() {
updates.insert("summary".to_string(), json!(text));
}
}
}
updates
});
// Build the graph with parallel execution
let agent = GraphAgent::builder("text_processor")
.description("Processes text with translation and summarization in parallel")
.channels(&["input", "translation", "summary", "result"])
.node(translator_node)
.node(summarizer_node)
.node_fn("combine", |ctx| async move {
let translation = ctx.get("translation").and_then(|v| v.as_str()).unwrap_or("N/A");
let summary = ctx.get("summary").and_then(|v| v.as_str()).unwrap_or("N/A");
let result = format!(
"=== Processing Complete ===\n\n\
French Translation:\n{}\n\n\
Summary:\n{}",
translation, summary
);
Ok(NodeOutput::new().with_update("result", json!(result)))
})
// Parallel execution: both nodes start simultaneously
.edge(START, "translator")
.edge(START, "summarizer")
.edge("translator", "combine")
.edge("summarizer", "combine")
.edge("combine", END)
.build()?;
// Execute the graph
let mut input = State::new();
input.insert("input".to_string(), json!("AI is transforming how we work and live."));
let result = agent.invoke(input, ExecutionConfig::new("thread-1")).await?;
println!("{}", result.get("result").and_then(|v| v.as_str()).unwrap_or(""));
Ok(())
}
SaΓda do exemplo:
=== Processing Complete ===
French Translation:
L'IA transforme notre faΓ§on de travailler et de vivre.
Summary:
AI is revolutionizing work and daily life through technological transformation.
Como funciona a execuΓ§Γ£o do grafo
VisΓ£o geral
Os agentes de grafo executam em superetapas: todos os nΓ³s prontos sΓ£o executados em paralelo, e entΓ£o o grafo espera que todos sejam concluΓdos antes da prΓ³xima etapa:
Step 1: START βββ¬βββΆ translator (running)
ββββΆ summarizer (running)
β³ Wait for both to complete...
Step 2: translator βββ¬βββΆ combine (running)
summarizer βββ
β³ Wait for combine to complete...
Step 3: combine βββΆ END β
Fluxo do estado pelos nΓ³s
Cada nΓ³ pode ler e gravar no estado compartilhado:
βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
β STEP 1: Initial state β
βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ€
β β
β State: { "input": "AI is transforming how we work" } β
β β
β β β
β β
β ββββββββββββββββββββ ββββββββββββββββββββ β
β β translator β β summarizer β β
β β reads "input" β β reads "input" β β
β ββββββββββββββββββββ ββββββββββββββββββββ β
β β
βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
β
βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
β STEP 2: After parallel execution β
βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ€
β β
β State: { β
β "input": "AI is transforming how we work", β
β "translation": "L'IA transforme notre faΓ§on de travailler", β
β "summary": "AI is revolutionizing work through technology" β
β } β
β β
β β β
β β
β ββββββββββββββββββββββββββββββββββββββββ β
β β combine β β
β β reads "translation" + "summary" β β
β β writes "result" β β
β ββββββββββββββββββββββββββββββββββββββββ β
β β
βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
β
βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
β STEP 3: Final state β
βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ€
β β
β State: { β
β "input": "AI is transforming how we work", β
β "translation": "L'IA transforme notre faΓ§on de travailler", β
β "summary": "AI is revolutionizing work through technology", β
β "result": "=== Processing Complete ===\n\nFrench..." β
β } β
β β
βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
O que faz isso funcionar
| Componente | FunΓ§Γ£o |
|---|---|
AgentNode | Envolve agentes LLM com mapeadores de entrada/saΓda |
input_mapper | Transforma o estado β entrada do agente Content |
output_mapper | Transforma eventos do agente β atualizaΓ§Γ΅es de estado |
channels | Declara os campos de estado que o grafo usarΓ‘ |
edge() | Define o fluxo de execuΓ§Γ£o entre os nΓ³s |
ExecutionConfig | Fornece o ID da thread para criar pontos de verificaΓ§Γ£o |
Roteamento condicional com classificaΓ§Γ£o de LLM
Crie sistemas de roteamento inteligentes nos quais LLMs decidam o caminho de execuΓ§Γ£o:
VisualizaΓ§Γ£o: roteamento baseado em sentimento
βββββββββββββββββββββββ
User Feedback β β
βββββββββββββββββΆ β CLASSIFIER β
β π§ Analyze tone β
ββββββββββββ¬βββββββββββ
β
βββββββββββββββββΌββββββββββββββββ
β β β
βΌ βΌ βΌ
ββββββββββββββββββββ ββββββββββββββββββββ ββββββββββββββββββββ
β POSITIVE β β NEGATIVE β β NEUTRAL β
β β β β β β
β π Thank you! β β π Apologize β β π Ask more β
β Celebrate β β Help fix β β questions β
ββββββββββββββββββββ ββββββββββββββββββββ ββββββββββββββββββββ
CΓ³digo completo do exemplo
use adk_agent::LlmAgentBuilder;
use adk_graph::{
edge::{END, Router, START},
graph::StateGraph,
node::{AgentNode, ExecutionConfig},
state::State,
};
use adk_model::GeminiModel;
use serde_json::json;
use std::sync::Arc;
#[tokio::main]
async fn main() -> anyhow::Result<()> {
dotenvy::dotenv().ok();
let api_key = std::env::var("GOOGLE_API_KEY")?;
let model = Arc::new(GeminiModel::new(&api_key, "gemini-3.7-flash")?);
// Create classifier agent
let classifier_agent = Arc::new(
LlmAgentBuilder::new("classifier")
.description("Classifies text sentiment")
.model(model.clone())
.instruction(
"You are a sentiment classifier. Analyze the input text and respond with \
ONLY one word: 'positive', 'negative', or 'neutral'. Nothing else.",
)
.build()?,
);
// Create response agents for each sentiment
let positive_agent = Arc::new(
LlmAgentBuilder::new("positive")
.description("Handles positive feedback")
.model(model.clone())
.instruction(
"You are a customer success specialist. The customer has positive feedback. \
Express gratitude, reinforce the positive experience, and suggest ways to \
share their experience. Be warm and appreciative. Keep response under 3 sentences.",
)
.build()?,
);
let negative_agent = Arc::new(
LlmAgentBuilder::new("negative")
.description("Handles negative feedback")
.model(model.clone())
.instruction(
"You are a customer support specialist. The customer has a complaint. \
Acknowledge their frustration, apologize sincerely, and offer help. \
Be empathetic. Keep response under 3 sentences.",
)
.build()?,
);
let neutral_agent = Arc::new(
LlmAgentBuilder::new("neutral")
.description("Handles neutral feedback")
.model(model.clone())
.instruction(
"You are a customer service representative. The customer has neutral feedback. \
Ask clarifying questions to better understand their needs. Be helpful and curious. \
Keep response under 3 sentences.",
)
.build()?,
);
// Create AgentNodes with mappers
let classifier_node = AgentNode::new(classifier_agent)
.with_input_mapper(|state| {
let text = state.get("feedback").and_then(|v| v.as_str()).unwrap_or("");
adk_core::Content::new("user").with_text(text)
})
.with_output_mapper(|events| {
let mut updates = std::collections::HashMap::new();
for event in events {
if let Some(content) = event.content() {
let text: String = content.parts.iter()
.filter_map(|p| p.text())
.collect::<Vec<_>>()
.join("")
.to_lowercase()
.trim()
.to_string();
let sentiment = if text.contains("positive") { "positive" }
else if text.contains("negative") { "negative" }
else { "neutral" };
updates.insert("sentiment".to_string(), json!(sentiment));
}
}
updates
});
// Response nodes (similar pattern for each)
let positive_node = AgentNode::new(positive_agent)
.with_input_mapper(|state| {
let text = state.get("feedback").and_then(|v| v.as_str()).unwrap_or("");
adk_core::Content::new("user").with_text(text)
})
.with_output_mapper(|events| {
let mut updates = std::collections::HashMap::new();
for event in events {
if let Some(content) = event.content() {
let text: String = content.parts.iter()
.filter_map(|p| p.text())
.collect::<Vec<_>>()
.join("");
updates.insert("response".to_string(), json!(text));
}
}
updates
});
// Build graph with conditional routing
let graph = StateGraph::with_channels(&["feedback", "sentiment", "response"])
.add_node(classifier_node)
.add_node(positive_node)
// ... add negative_node and neutral_node similarly
.add_edge(START, "classifier")
.add_conditional_edges(
"classifier",
Router::by_field("sentiment"), // Route based on sentiment field
[
("positive", "positive"),
("negative", "negative"),
("neutral", "neutral"),
],
)
.add_edge("positive", END)
.add_edge("negative", END)
.add_edge("neutral", END)
.compile()?;
// Test with different feedback
let mut input = State::new();
input.insert("feedback".to_string(), json!("Your product is amazing! I love it!"));
let result = graph.invoke(input, ExecutionConfig::new("feedback-1")).await?;
println!("Sentiment: {}", result.get("sentiment").and_then(|v| v.as_str()).unwrap_or(""));
println!("Response: {}", result.get("response").and_then(|v| v.as_str()).unwrap_or(""));
Ok(())
}
Fluxo do exemplo:
Input: "Your product is amazing! I love it!"
β
Classifier: "positive"
β
Positive Agent: "Thank you so much for the wonderful feedback!
We're thrilled you love our product.
Would you consider leaving a review to help others?"
PadrΓ£o ReAct: raciocΓnio + aΓ§Γ£o
Crie agentes que possam usar ferramentas iterativamente para resolver problemas complexos:
VisualizaΓ§Γ£o: ciclo ReAct
βββββββββββββββββββββββ
User Question β β
βββββββββββββββββΆ β REASONER β
β π§ Think + Act β
ββββββββββββ¬βββββββββββ
β
βΌ
βββββββββββββββββββββββ
β Has tool calls? β
β β
ββββββββββββ¬βββββββββββ
β
βββββββββββββββββ΄ββββββββββββββββ
β β
βΌ βΌ
ββββββββββββββββββββ ββββββββββββββββββββ
β YES β β NO β
β β β β
β π Loop back β β β
Final answer β
β to reasoner β β END β
βββββββββββ¬βββββββββ ββββββββββββββββββββ
β
βββββββββββββββββββ
β
βΌ
βββββββββββββββββββββββ
β REASONER β
β π§ Think + Act β
β (next iteration) β
βββββββββββββββββββββββ
Exemplo completo de ReAct
use adk_agent::LlmAgentBuilder;
use adk_core::{Part, Tool};
use adk_graph::{
edge::{END, START},
graph::StateGraph,
node::{AgentNode, ExecutionConfig, NodeOutput},
state::State,
};
use adk_model::GeminiModel;
use adk_tool::FunctionTool;
use serde_json::json;
use std::sync::Arc;
#[tokio::main]
async fn main() -> anyhow::Result<()> {
dotenvy::dotenv().ok();
let api_key = std::env::var("GOOGLE_API_KEY")?;
let model = Arc::new(GeminiModel::new(&api_key, "gemini-3.7-flash")?);
// Create tools
let weather_tool = Arc::new(FunctionTool::new(
"get_weather",
"Get the current weather for a location. Takes a 'location' parameter (city name).",
|_ctx, args| async move {
let location = args.get("location").and_then(|v| v.as_str()).unwrap_or("unknown");
Ok(json!({
"location": location,
"temperature": "72Β°F",
"condition": "Sunny",
"humidity": "45%"
}))
},
)) as Arc<dyn Tool>;
let calculator_tool = Arc::new(FunctionTool::new(
"calculator",
"Perform mathematical calculations. Takes an 'expression' parameter (string).",
|_ctx, args| async move {
let expr = args.get("expression").and_then(|v| v.as_str()).unwrap_or("0");
let result = match expr {
"2 + 2" => "4",
"10 * 5" => "50",
"100 / 4" => "25",
"15 - 7" => "8",
_ => "Unable to evaluate",
};
Ok(json!({ "result": result, "expression": expr }))
},
)) as Arc<dyn Tool>;
// Create reasoner agent with tools
let reasoner_agent = Arc::new(
LlmAgentBuilder::new("reasoner")
.description("Reasoning agent with tools")
.model(model.clone())
.instruction(
"You are a helpful assistant with access to tools. Use tools when needed to answer questions. \
When you have enough information, provide a final answer without using more tools.",
)
.tool(weather_tool)
.tool(calculator_tool)
.build()?,
);
// Create reasoner node that detects tool usage
let reasoner_node = AgentNode::new(reasoner_agent)
.with_input_mapper(|state| {
let question = state.get("question").and_then(|v| v.as_str()).unwrap_or("");
adk_core::Content::new("user").with_text(question)
})
.with_output_mapper(|events| {
let mut updates = std::collections::HashMap::new();
let mut has_tool_calls = false;
let mut response = String::new();
for event in events {
if let Some(content) = event.content() {
for part in &content.parts {
match part {
Part::FunctionCall { .. } => {
has_tool_calls = true;
}
Part::Text { text } => {
response.push_str(text);
}
_ => {}
}
}
}
}
updates.insert("has_tool_calls".to_string(), json!(has_tool_calls));
updates.insert("response".to_string(), json!(response));
updates
});
// Build ReAct graph with cycle
let graph = StateGraph::with_channels(&["question", "has_tool_calls", "response", "iteration"])
.add_node(reasoner_node)
.add_node_fn("counter", |ctx| async move {
let i = ctx.get("iteration").and_then(|v| v.as_i64()).unwrap_or(0);
Ok(NodeOutput::new().with_update("iteration", json!(i + 1)))
})
.add_edge(START, "counter")
.add_edge("counter", "reasoner")
.add_conditional_edges(
"reasoner",
|state| {
let has_tools = state.get("has_tool_calls").and_then(|v| v.as_bool()).unwrap_or(false);
let iteration = state.get("iteration").and_then(|v| v.as_i64()).unwrap_or(0);
// Safety limit
if iteration >= 5 { return END.to_string(); }
if has_tools {
"counter".to_string() // Loop back for more reasoning
} else {
END.to_string() // Done - final answer
}
},
[("counter", "counter"), (END, END)],
)
.compile()?
.with_recursion_limit(10);
// Test the ReAct agent
let mut input = State::new();
input.insert("question".to_string(), json!("What's the weather in Paris and what's 15 + 25?"));
let result = graph.invoke(input, ExecutionConfig::new("react-1")).await?;
println!("Final answer: {}", result.get("response").and_then(|v| v.as_str()).unwrap_or(""));
println!("Iterations: {}", result.get("iteration").and_then(|v| v.as_i64()).unwrap_or(0));
Ok(())
}
Fluxo do exemplo:
Question: "What's the weather in Paris and what's 15 + 25?"
Iteration 1:
Reasoner: "I need to get weather info and do math"
β Calls get_weather(location="Paris") and calculator(expression="15 + 25")
β has_tool_calls = true β Loop back
Iteration 2:
Reasoner: "Based on the results: Paris is 72Β°F and sunny, 15 + 25 = 40"
β No tool calls β has_tool_calls = false β END
Final Answer: "The weather in Paris is 72Β°F and sunny with 45% humidity.
And 15 + 25 equals 40."
AgentNode
Envolve qualquer ADK Agent (normalmente LlmAgent) como um nΓ³ do grafo:
O que o agente vΓͺ
Um agente dentro de um grafo Γ© executado em um contexto derivado da invocaΓ§Γ£o que
iniciou o grafo, portanto, comporta-se da mesma forma que fora dele. Quando um Runner
invoca um GraphAgent, a identidade e os serviΓ§os do chamador sΓ£o transmitidos
automaticamente:
| Transportado por | ObservaΓ§Γ£o |
|---|---|
app_name, user_id, session_id | Do chamador, nΓ£o um sintΓ©tico |
| Escopos e metadados da solicitaΓ§Γ£o | Para que as verificaΓ§Γ΅es de escopo vejam as concessΓ΅es do chamador |
| ServiΓ§o de segredos, memΓ³ria, artefatos, estado compartilhado | DisponΓveis exatamente como fora do grafo |
| Cancelamento | Runner::interrupt alcanΓ§a um agente em execuΓ§Γ£o como um nΓ³ |
RunConfig | Herdado do chamador |
branch | Derivado, como {caller_branch}.{agent_name}, para que os eventos de um nΓ³ possam ser atribuΓdos |
Um grafo invocado diretamente β graph.invoke(state, ExecutionConfig::new("thread")) β
nΓ£o tem nenhuma invocaΓ§Γ£o para herdar. Esse Γ© o modo independente: o nΓ³ recebe
user_id = "graph_user", app_name = "graph_app", a ramificaΓ§Γ£o main, nenhum segredo e nenhuma
memΓ³ria. Γ um modo deliberado para executar um grafo fora de um Runner, nΓ£o um
recurso alternativo a ser usado em produΓ§Γ£o.
Para fazer a ponte manualmente β por exemplo, ao controlar um grafo a partir do seu prΓ³prio executor β passe a invocaΓ§Γ£o explicitamente:
let config = ExecutionConfig::new(ctx.session_id()).with_parent_context(ctx.clone());
ObservaΓ§Γ£o: o nΓ³ ainda Γ© executado em sua prΓ³pria sessΓ£o de grafo na memΓ³ria, portanto o histΓ³rico de conversas do agente dentro de um nΓ³ fica restrito ao nΓ³, em vez de ser acrescentado Γ sessΓ£o do chamador.
let node = AgentNode::new(llm_agent)
.with_input_mapper(|state| {
// Transform graph state to agent input Content
let text = state.get("input").and_then(|v| v.as_str()).unwrap_or("");
adk_core::Content::new("user").with_text(text)
})
.with_output_mapper(|events| {
// Transform agent events to state updates
let mut updates = std::collections::HashMap::new();
for event in events {
if let Some(content) = event.content() {
let text: String = content.parts.iter()
.filter_map(|p| p.text())
.collect::<Vec<_>>()
.join("");
updates.insert("output".to_string(), json!(text));
}
}
updates
});
NΓ³s de funΓ§Γ£o
FunΓ§Γ΅es assΓncronas simples que processam o estado:
.node_fn("process", |ctx| async move {
let input = ctx.state.get("input").unwrap();
let output = process_data(input).await?;
Ok(NodeOutput::new().with_update("output", output))
})
Tipos de aresta
Arestas estΓ‘ticas
ConexΓ΅es diretas entre nΓ³s:
.edge(START, "first_node")
.edge("first_node", "second_node")
.edge("second_node", END)
Arestas condicionais
Roteamento dinΓ’mico com base no estado:
.conditional_edge(
"router",
|state| {
match state.get("next").and_then(|v| v.as_str()) {
Some("research") => "research_node".to_string(),
Some("write") => "write_node".to_string(),
_ => END.to_string(),
}
},
[
("research_node", "research_node"),
("write_node", "write_node"),
(END, END),
],
)
Roteamento a partir de dentro de um nΓ³
Uma aresta condicional define seus destinos quando o grafo Γ© construΓdo. NodeOutput::with_goto
nΓ£o faz isso: um nΓ³ grava o estado e nomeia seus sucessores na mesma etapa, e pode
nomear qualquer nΓ³ do grafo, inclusive um ao qual nΓ£o tenha nenhuma aresta.
use adk_graph::node::NodeOutput;
use serde_json::json;
// The node decides where control goes, from what it just computed.
async fn triage(ctx: &adk_graph::node::NodeContext) -> adk_graph::error::Result<NodeOutput> {
let amount = ctx.get("amount").and_then(|v| v.as_f64()).unwrap_or(0.0);
let next = if amount > 10_000.0 { "escalate" } else { "auto_approve" };
Ok(NodeOutput::new().with_update("risk", json!(next)).with_goto([next]))
}
| Comportamento | Regra |
|---|---|
| Arestas declaradas | Um nΓ³ que define um goto tambΓ©m nΓ£o segue suas arestas de saΓda. O goto as substitui. |
| VΓ‘rios destinos | Todos os nΓ³s nomeados sΓ£o executados, admitidos em ordem classificada. |
END | Nomear END interrompe esse ramo. |
| Um nome desconhecido | A execuΓ§Γ£o falha com GraphError::UnknownRouteTarget. |
| Sem goto | As arestas declaradas decidem, que Γ© o padrΓ£o. |
A fronteira que um goto produz recebe checkpoint como qualquer outra, portanto uma execuΓ§Γ£o pausada Γ© retomada no nΓ³ escolhido pelo goto.
Nota: use
add_conditional_edgesquando os possΓveis destinos forem conhecidos ao criar o grafo β as arestas aparecerΓ£o entΓ£o em um diagrama renderizado. Use um goto quando a escolha pertencer ao nΓ³.
Auxiliares de roteamento
Use roteadores integrados para padrΓ΅es comuns:
use adk_graph::edge::Router;
// Route based on a state field value
.conditional_edge("classifier", Router::by_field("sentiment"), [
("positive", "positive_handler"),
("negative", "negative_handler"),
("neutral", "neutral_handler"),
])
// Route based on boolean field
.conditional_edge("check", Router::by_bool("approved"), [
("true", "execute"),
("false", "reject"),
])
// Limit iterations
.conditional_edge("loop", Router::max_iterations("count", 5), [
("continue", "process"),
("done", END),
])
ExecuΓ§Γ£o paralela
MΓΊltiplas arestas a partir de um ΓΊnico nΓ³ sΓ£o executadas em paralelo:
let agent = GraphAgent::builder("parallel_processor")
.channels(&["input", "translation", "summary", "analysis"])
.node(translator_node)
.node(summarizer_node)
.node(analyzer_node)
.node(combiner_node)
// All three start simultaneously
.edge(START, "translator")
.edge(START, "summarizer")
.edge(START, "analyzer")
// Wait for all to complete before combining
.edge("translator", "combiner")
.edge("summarizer", "combiner")
.edge("analyzer", "combiner")
.edge("combiner", END)
.build()?;
combiner Γ© executado uma vez, depois que todas as trΓͺs ramificaΓ§Γ΅es chegam. Nada no cΓ³digo acima solicita isso: um nΓ³ com mais de uma aresta direta de entrada Γ© automaticamente adiado no momento da compilaΓ§Γ£o. Portanto, ramificaΓ§Γ΅es de comprimentos diferentes sΓ£o unidas corretamente sem configuraΓ§Γ£o.
Dois detalhes decorrem da forma como a contagem Γ© realizada:
| Caso | Comportamento |
|---|---|
| Predecessores condicionais | NΓ£o sΓ£o contabilizados. Um ramo condicional pode nunca ser executado, portanto esperar por ele poderia bloquear a junΓ§Γ£o. |
| Um quΓ³rum em vez de todos | Defina min_predecessors em DeferredNodeConfig e marque o nΓ³ com mark_deferred para liberar apΓ³s a chegada de n de m. |
Grafos CΓclicos (PadrΓ£o ReAct)
Crie agentes de raciocΓnio iterativo com ciclos:
use adk_core::Part;
// Create agent with tools
let reasoner = Arc::new(
LlmAgentBuilder::new("reasoner")
.model(model)
.instruction("Use tools to answer questions. Provide final answer when done.")
.tool(search_tool)
.tool(calculator_tool)
.build()?
);
let reasoner_node = AgentNode::new(reasoner)
.with_input_mapper(|state| {
let question = state.get("question").and_then(|v| v.as_str()).unwrap_or("");
adk_core::Content::new("user").with_text(question)
})
.with_output_mapper(|events| {
let mut updates = std::collections::HashMap::new();
let mut has_tool_calls = false;
let mut response = String::new();
for event in events {
if let Some(content) = event.content() {
for part in &content.parts {
match part {
Part::FunctionCall { name, .. } => {
has_tool_calls = true;
}
Part::Text { text } => {
response.push_str(text);
}
_ => {}
}
}
}
}
updates.insert("has_tool_calls".to_string(), json!(has_tool_calls));
updates.insert("response".to_string(), json!(response));
updates
});
// Build graph with cycle
let react_agent = StateGraph::with_channels(&["question", "has_tool_calls", "response", "iteration"])
.add_node(reasoner_node)
.add_node_fn("counter", |ctx| async move {
let i = ctx.get("iteration").and_then(|v| v.as_i64()).unwrap_or(0);
Ok(NodeOutput::new().with_update("iteration", json!(i + 1)))
})
.add_edge(START, "counter")
.add_edge("counter", "reasoner")
.add_conditional_edges(
"reasoner",
|state| {
let has_tools = state.get("has_tool_calls").and_then(|v| v.as_bool()).unwrap_or(false);
let iteration = state.get("iteration").and_then(|v| v.as_i64()).unwrap_or(0);
// Safety limit
if iteration >= 5 { return END.to_string(); }
if has_tools {
"counter".to_string() // Loop back
} else {
END.to_string() // Done
}
},
[("counter", "counter"), (END, END)],
)
.compile()?
.with_recursion_limit(10);
Supervisor Multiagente
Encaminhe tarefas para agentes especialistas:
// Create supervisor agent
let supervisor = Arc::new(
LlmAgentBuilder::new("supervisor")
.model(model.clone())
.instruction("Route tasks to: researcher, writer, or coder. Reply with agent name only.")
.build()?
);
let supervisor_node = AgentNode::new(supervisor)
.with_output_mapper(|events| {
let mut updates = std::collections::HashMap::new();
for event in events {
if let Some(content) = event.content() {
let text: String = content.parts.iter()
.filter_map(|p| p.text())
.collect::<Vec<_>>()
.join("")
.to_lowercase();
let next = if text.contains("researcher") { "researcher" }
else if text.contains("writer") { "writer" }
else if text.contains("coder") { "coder" }
else { "done" };
updates.insert("next_agent".to_string(), json!(next));
}
}
updates
});
// Build supervisor graph
let graph = StateGraph::with_channels(&["task", "next_agent", "research", "content", "code"])
.add_node(supervisor_node)
.add_node(researcher_node)
.add_node(writer_node)
.add_node(coder_node)
.add_edge(START, "supervisor")
.add_conditional_edges(
"supervisor",
Router::by_field("next_agent"),
[
("researcher", "researcher"),
("writer", "writer"),
("coder", "coder"),
("done", END),
],
)
// Agents report back to supervisor
.add_edge("researcher", "supervisor")
.add_edge("writer", "supervisor")
.add_edge("coder", "supervisor")
.compile()?;
Gerenciamento de Estado
Esquema de Estado com Redutores
Controle como as atualizaΓ§Γ΅es de estado sΓ£o mescladas:
let schema = StateSchema::builder()
.channel("current_step") // Overwrite (default)
.list_channel("messages") // Append to list
.channel_with_reducer("count", Reducer::Sum) // Sum values
.channel_with_reducer("data", Reducer::Custom(Arc::new(|old, new| {
// Custom merge logic
merge_json(old, new)
})))
.build();
let agent = GraphAgent::builder("stateful")
.state_schema(schema)
// ... nodes and edges
.build()?;
Tipos de Redutor
| Redutor | Comportamento |
|---|---|
Overwrite | Substitui o valor antigo pelo novo (padrΓ£o) |
Append | Anexa Γ lista |
Sum | Adicionar valores numΓ©ricos |
Custom | FunΓ§Γ£o de mesclagem personalizada |
Ordem de atualizaΓ§Γ£o
Os nΓ³s em uma ΓΊnica superetapa sΓ£o executados simultaneamente e terminam na ordem que o trabalho deles exigir. As atualizaΓ§Γ΅es de estado sΓ£o aplicadas na ordem do nome do nΓ³, nΓ£o na ordem em que os nΓ³s terminaram.
A ordem importa sempre que um redutor nΓ£o Γ© comutativo. Append cria um array, portanto a ordem Γ© o resultado; um redutor Custom tambΓ©m pode ser sensΓvel Γ ordem. A ordenaΓ§Γ£o pelo nome do nΓ³ torna uma execuΓ§Γ£o reproduzΓvel: o mesmo grafo e a mesma entrada produzem o mesmo estado, independentemente do momento em que uma dependΓͺncia lenta termina.
Quando vΓ‘rios canais sΓ£o escritos por um nΓ³, eles sΓ£o aplicados na ordem do nome do canal. Os canais nΓ£o interagem, portanto isso importa apenas para a leitura de um rastreamento.
CriaΓ§Γ£o de pontos de verificaΓ§Γ£o
Habilite o estado persistente para tolerΓ’ncia a falhas e participaΓ§Γ£o humana no circuito:
Em memΓ³ria (desenvolvimento)
use adk_graph::checkpoint::MemoryCheckpointer;
let checkpointer = Arc::new(MemoryCheckpointer::new());
let graph = StateGraph::with_channels(&["task", "result"])
// ... nodes and edges
.compile()?
.with_checkpointer_arc(checkpointer.clone());
SQLite (produΓ§Γ£o)
use adk_graph::checkpoint::SqliteCheckpointer;
let checkpointer = SqliteCheckpointer::new("checkpoints.db").await?;
let graph = StateGraph::with_channels(&["task", "result"])
// ... nodes and edges
.compile()?
.with_checkpointer(checkpointer);
O que um ponto de verificaΓ§Γ£o registra
Um ponto de verificaΓ§Γ£o armazena o estado acumulado, o nΓΊmero da etapa e a fronteira β os nΓ³s que ainda precisam ser executados. Ele Γ© gravado depois que a fronteira avanΓ§a, portanto a retomada nunca executa novamente um nΓ³ que jΓ‘ foi concluΓdo e nunca aplica suas atualizaΓ§Γ΅es duas vezes. Uma execuΓ§Γ£o que termina cria um ponto de verificaΓ§Γ£o com uma fronteira vazia, portanto retomar uma thread concluΓda retorna o estado final em vez de reiniciar o grafo.
Dois casos criam deliberadamente um ponto de verificaΓ§Γ£o da fronteira que estava em execuΓ§Γ£o, em vez da prΓ³xima, porque o nΓ³ interrompido ainda nΓ£o produziu suas atualizaΓ§Γ΅es e precisa ser executado novamente ao retomar:
| SituaΓ§Γ£o | Fronteira salva |
|---|---|
| Superetapa concluΓda | Os prΓ³ximos nΓ³s a executar |
| ExecuΓ§Γ£o concluΓda | Vazia |
| InterrupΓ§Γ£o gerada (bloqueante ou em streaming) | Os nΓ³s que estavam em execuΓ§Γ£o |
As execuΓ§Γ΅es em streaming criam pontos de verificaΓ§Γ£o na mesma programaΓ§Γ£o que as execuΓ§Γ΅es bloqueantes, inclusive quando uma interrupΓ§Γ£o encerra o stream, portanto uma pausa com participaΓ§Γ£o humana pode ser retomada em qualquer um dos modos de execuΓ§Γ£o.
HistΓ³rico de Pontos de VerificaΓ§Γ£o (Viagem no Tempo)
Somente leitura.
TimeTravelHandle::state_history(from, to)retorna o estado que foi armazenado em cada etapa com ponto de verificaΓ§Γ£o. Ele nΓ£o executa nada β nenhum nΓ³ Γ© executado, nenhum evento Γ© regenerado e nenhum efeito colateral se repete. Para executar novamente a partir de um ponto do histΓ³rico, usefork_atpara criar um branch desse ponto de verificaΓ§Γ£o e invoque o grafo na thread derivada. O mΓ©todo anteriormente se chamavareplaye era documentado como reexecutando o grafo, o que ele nunca fez.
Os pontos de verificaΓ§Γ£o tambΓ©m permitem a retomada durΓ‘vel β se a execuΓ§Γ£o de um grafo falhar ou o processo for reiniciado, a execuΓ§Γ£o serΓ‘ retomada a partir do ΓΊltimo ponto de verificaΓ§Γ£o persistido, em vez de comeΓ§ar novamente. Use SqliteCheckpointer (o recurso sqlite) para obter persistΓͺncia segura contra falhas. MemoryCheckpointer mantΓ©m os pontos de verificaΓ§Γ£o no processo, portanto eles nΓ£o sobrevivem a uma reinicializaΓ§Γ£o. Esses sΓ£o os dois backends fornecidos por este crate; implemente a trait Checkpointer para qualquer outra opΓ§Γ£o.
// List all checkpoints for a thread
let checkpoints = checkpointer.list("thread-id").await?;
for cp in checkpoints {
println!("Step {}: {:?}", cp.step, cp.state.get("status"));
}
// Load a specific checkpoint
if let Some(checkpoint) = checkpointer.load_by_id(&checkpoint_id).await? {
println!("State at step {}: {:?}", checkpoint.step, checkpoint.state);
}
ParticipaΓ§Γ£o Humana
Pause a execuΓ§Γ£o para obter aprovaΓ§Γ£o humana usando interrupΓ§Γ΅es dinΓ’micas:
use adk_graph::{error::GraphError, node::NodeOutput};
// Planner agent assesses risk
let planner_node = AgentNode::new(planner_agent)
.with_output_mapper(|events| {
let mut updates = std::collections::HashMap::new();
for event in events {
if let Some(content) = event.content() {
let text: String = content.parts.iter()
.filter_map(|p| p.text())
.collect::<Vec<_>>()
.join("");
// Extract risk level from LLM response
let risk = if text.to_lowercase().contains("risk: high") { "high" }
else if text.to_lowercase().contains("risk: medium") { "medium" }
else { "low" };
updates.insert("plan".to_string(), json!(text));
updates.insert("risk_level".to_string(), json!(risk));
}
}
updates
});
// Review node with dynamic interrupt
let graph = StateGraph::with_channels(&["task", "plan", "risk_level", "approved", "result"])
.add_node(planner_node)
.add_node(executor_node)
.add_node_fn("review", |ctx| async move {
let risk = ctx.get("risk_level").and_then(|v| v.as_str()).unwrap_or("low");
let approved = ctx.get("approved").and_then(|v| v.as_bool());
// Already approved - continue
if approved == Some(true) {
return Ok(NodeOutput::new());
}
// High/medium risk - interrupt for approval
if risk == "high" || risk == "medium" {
return Ok(NodeOutput::interrupt_with_data(
&format!("{} RISK: Human approval required", risk.to_uppercase()),
json!({
"plan": ctx.get("plan"),
"risk_level": risk,
"action": "Set 'approved' to true to continue"
})
));
}
// Low risk - auto-approve
Ok(NodeOutput::new().with_update("approved", json!(true)))
})
.add_edge(START, "planner")
.add_edge("planner", "review")
.add_edge("review", "executor")
.add_edge("executor", END)
.compile()?
.with_checkpointer_arc(checkpointer.clone());
// Execute - may pause for approval
let thread_id = "task-001";
let result = graph.invoke(input, ExecutionConfig::new(thread_id)).await;
match result {
Err(GraphError::Interrupted(interrupt)) => {
println!("*** EXECUTION PAUSED ***");
println!("Reason: {}", interrupt.interrupt);
println!("Plan awaiting approval: {:?}", interrupt.state.get("plan"));
// Human reviews and approves...
// Update state with approval
graph.update_state(thread_id, [("approved".to_string(), json!(true))]).await?;
// Resume execution
let final_result = graph.invoke(State::new(), ExecutionConfig::new(thread_id)).await?;
println!("Final result: {:?}", final_result.get("result"));
}
Ok(result) => {
println!("Completed without interrupt: {:?}", result);
}
Err(e) => {
println!("Error: {}", e);
}
}
InterrupΓ§Γ΅es EstΓ‘ticas
Use interrupt_before ou interrupt_after para pontos de pausa obrigatΓ³rios:
let graph = StateGraph::with_channels(&["task", "plan", "result"])
.add_node(planner_node)
.add_node(executor_node)
.add_edge(START, "planner")
.add_edge("planner", "executor")
.add_edge("executor", END)
.compile()?
.with_interrupt_before(&["executor"]); // Always pause before execution
ExecuΓ§Γ£o em Streaming
Transmita eventos Γ medida que o grafo Γ© executado:
use futures::StreamExt;
use adk_graph::stream::StreamMode;
let stream = agent.stream(input, config, StreamMode::Updates);
while let Some(event) = stream.next().await {
match event? {
StreamEvent::NodeStart(name) => println!("Starting: {}", name),
StreamEvent::Updates { node, updates } => {
println!("{} updated state: {:?}", node, updates);
}
StreamEvent::NodeEnd(name) => println!("Completed: {}", name),
StreamEvent::Done(state) => println!("Final state: {:?}", state),
_ => {}
}
}
Modos de Streaming
| Modo | DescriΓ§Γ£o |
|---|---|
Values | Transmite o estado completo apΓ³s cada nΓ³ |
Updates | Transmite apenas as alteraΓ§Γ΅es de estado |
Messages | Transmitir atualizaΓ§Γ΅es de tipo de mensagem |
Debug | Transmitir todos os eventos internos |
O modo Messages lΓͺ tokens de Node::execute_stream Γ medida que sΓ£o produzidos.
Cada nΓ³ Γ© executado uma vez por superetapa nesse modo: o nΓ³ relata suas
atualizaΓ§Γ΅es de estado no fluxo como um evento StreamEvent::Updates, e o executor aplica
essas atualizaΓ§Γ΅es em vez de executar o nΓ³ uma segunda vez para coletΓ‘-las. Isso Γ©
mais importante para AgentNode, em que uma segunda execuΓ§Γ£o significaria uma
segunda chamada faturada ao modelo por nΓ³.
Importante: um
Nodepersonalizado que substituiexecute_streamdeve produzir um eventoStreamEvent::Updatescontendo suas atualizaΓ§Γ΅es de estado. Sem isso, o nΓ³ transmite eventos, mas nΓ£o contribui com estado no modoMessages. Oexecute_streampadrΓ£o, que encapsulaexecute, faz isso por vocΓͺ.
As polΓticas de tempo limite se aplicam Γ prΓ³pria execuΓ§Γ£o transmitida. Para um fluxo,
idle_timeout significa que nenhum evento foi produzido dentro do limite.
Subgrafos
Um grafo compilado Γ© executado como um nΓ³ de outro por meio de SubgraphNode. O grafo
interno mantΓ©m seus prΓ³prios canais, arestas e portas de interrupΓ§Γ£o, e troca canais
nomeados com seu pai.
use adk_graph::subgraph::SubgraphNode;
use std::sync::Arc;
let outer = StateGraph::with_channels(&["document", "size"])
.add_node(
SubgraphNode::new("measure_doc", Arc::new(inner))
.with_input("document", "text")
.with_output("length", "size"),
)
.add_edge(START, "measure_doc")
.add_edge("measure_doc", END)
.compile()?;
| Regra | Comportamento |
|---|---|
| Nomes compartilhados | Um canal que ambos os esquemas declaram com o mesmo nome passa nos dois sentidos. |
isolated() | Nada passa implicitamente; toda troca deve ser nomeada. Vale a pena quando os dois grafos sΓ£o mantidos separadamente, pois adicionar um canal a um deles nΓ£o poderΓ‘ comeΓ§ar a alimentar silenciosamente o outro. |
| Uma pausa dentro | Pausa o pai, carregando o nome do subgrafo e a mensagem interna. |
| Threads | O subgrafo Γ© executado em <parent thread>/<node name>, portanto dois subgrafos de um mesmo pai nΓ£o podem colidir. |
| Um nome de canal incorreto | Falha quando o pai compila, identificando o canal e o lado. |
Essa ΓΊltima linha Γ© a diferenΓ§a que vale a pena conhecer: ambos os esquemas estΓ£o disponΓveis antes de qualquer execuΓ§Γ£o, portanto um mapeamento que nomeie um canal que nenhum dos lados declara nΓ£o pode chegar a uma execuΓ§Γ£o e aparecer como um valor ausente. Um subgrafo que nΓ£o troque absolutamente nada Γ© rejeitado da mesma forma, pois nΓ£o poderia afetar seu pai.
Retomando uma pausa dentro de um subgrafo
Nada mais Γ© necessΓ‘rio. Invocar o pai novamente na mesma thread entra novamente no subgrafo, que encontra seu prΓ³prio checkpoint em <parent thread>/<node name> e continua de onde parou. O trabalho concluΓdo pelo subgrafo antes da pausa nΓ£o Γ© repetido, e uma pausa ocorrida vΓ‘rios nΓveis abaixo Γ© retomada da mesma forma β a mensagem nomeia cada nΓvel pelo qual passou.
Um subgrafo que declara um ponto de interrupΓ§Γ£o, mas nΓ£o possui nenhum armazenador de checkpoints, Γ© rejeitado quando o pai Γ© compilado: ele entraria novamente em seu primeiro nΓ³ e pagaria uma segunda vez pelo trabalho jΓ‘ concluΓdo.
| Tipo de pausa | Como a resposta chega |
|---|---|
interrupt_before / interrupt_after dentro | Nada a fornecer; a retomada libera o bloqueio que foi acionado |
| Um nΓ³ interno decidindo por conta prΓ³pria | A decisΓ£o chega como estado, projetada por meio do mapeamento de canais |
Como ambos os grafos mantΓͺm checkpointers reais, isso sobrevive a uma reinicializaΓ§Γ£o do processo: um novo conjunto de objetos de grafo que compartilha apenas os bancos de dados retoma a mesma execuΓ§Γ£o.
Devolvendo o controle ao pai
Um nΓ³ dentro de um subgrafo pode encerrar o prΓ³prio grafo e nomear um nΓ³ do grafo que o contΓ©m:
Ok(NodeOutput::new()
.with_update("reason", json!("no confident answer"))
.with_goto_parent(["escalate"]))
O subgrafo termina e projeta seus canais de saΓda normalmente; em seguida, o pai continua em escalate em vez de seguir as arestas declaradas pelo nΓ³ do subgrafo. O pai valida o destino, pois somente ele conhece seus prΓ³prios nΓ³s.
Controles de confiabilidade e custo
Cada um destes recursos vem desativado por padrΓ£o, portanto um grafo se comporta como antes de vocΓͺ configurar qualquer um deles.
Nova tentativa por nΓ³
Uma falha transitΓ³ria β um limite de taxa ou uma conexΓ£o interrompida β caso contrΓ‘rio encerra a execuΓ§Γ£o.
use adk_graph::retry::{RetryOn, RetryPolicy};
use std::time::Duration;
let graph = graph.with_node_retry(
"call_model",
RetryPolicy::new(3)
.with_initial_delay(Duration::from_millis(500))
.with_max_delay(Duration::from_secs(8))
.with_backoff_factor(2.0)
.with_retry_on(RetryOn::Any),
);
O atraso aumenta em backoff_factor, Γ© limitado a max_delay e, em seguida, recebe jitter. Um nΓ³ sem polΓtica Γ© executado uma vez, portanto as novas tentativas continuam sendo opcionais. Uma polΓtica de RetryPolicy::default() permite dez tentativas, cujas nove esperas totalizam cerca de 243 segundos β reduza max_attempts quando um chamador estiver aguardando a resposta.
Uma interrupΓ§Γ£o nunca Γ© repetida, independentemente do que retry_on diga: uma pausa nΓ£o Γ© uma falha. A contagem de tentativas Γ© armazenada no checkpoint, portanto uma execuΓ§Γ£o retomada continua usando o orΓ§amento existente em vez de reiniciΓ‘-lo.
Limitando a concorrΓͺncia
Uma distribuiΓ§Γ£o ampla despacha toda a sua fronteira de uma sΓ³ vez, o que pode esgotar um pool de conexΓ΅es ou acionar um limite de taxa do provedor.
let graph = graph.with_max_concurrency(4);
Os nΓ³s alΓ©m do limite aguardam uma vaga. A ordem de admissΓ£o Γ© a fronteira classificada por nome, portanto nΓ£o depende do tempo. As invocaΓ§Γ΅es imperativas de filhos ficam fora desse orΓ§amento, pois um pai aguarda seus filhos enquanto mantΓ©m sua prΓ³pria vaga.
Tempos limite por nΓ³
Um TimeoutPolicy limita uma ΓΊnica tentativa e, com idle_timeout, define por quanto tempo um nΓ³ pode ficar sem relatar progresso. Exceder qualquer um deles gera GraphError::NodeTimedOut, sobre o qual uma polΓtica de novas tentativas pode entΓ£o agir.
Invocando um nΓ³ diretamente
Quando o nΓΊmero de subtarefas vem do estado, e nΓ£o da estrutura do grafo, um nΓ³ pode invocar outro nΓ³ por conta prΓ³pria:
use adk_graph::child::RunNodeOptions;
let output = ctx
.run_node_with("reviewer", json!({ "aspect": aspect }), RunNodeOptions::with_run_id(aspect))
.await?;
O destino nΓ£o precisa de uma aresta. Cada filho concluΓdo Γ© registrado em <parent>/<child>@<run_id>, portanto uma execuΓ§Γ£o retomada retorna a resposta registrada em vez de executar o filho novamente β o que Γ© importante quando o filho exige uma chamada de modelo.
Armazenamento em cache de nΓ³s
cache_policy em um nΓ³ identifica o resultado por nome do nΓ³ e estado atual, com um TTL opcional, portanto uma entrada inalterada ignora o trabalho. Requer o recurso node-cache; um armazenamento baseado em Redis estΓ‘ disponΓvel por meio de redis-cache.
Pontos de verificaΓ§Γ£o delta
O recurso delta armazena a diferenΓ§a entre superetapas em vez do estado inteiro, o que Γ© importante quando o estado Γ© grande e hΓ‘ muitas etapas.
Viagem no tempo
Com o recurso time-travel, graph.time_travel(thread_id)? retorna um identificador sobre o histΓ³rico de pontos de verificaΓ§Γ£o de uma thread: listar as etapas, ler o estado em uma delas ou fork_at um ponto de verificaΓ§Γ£o para criar uma nova thread a partir dele. A chamada retorna Result porque toda operaΓ§Γ£o lΓͺ pontos de verificaΓ§Γ£o; portanto, um grafo sem um registrador de pontos de verificaΓ§Γ£o informa GraphError::CheckpointError em vez de entrar em pΓ’nico.
PadrΓ΅es para todo o grafo
Repetir a mesma nova tentativa em vinte nΓ³s Γ© fΓ‘cil de fazer incorretamente por omissΓ£o.
use adk_graph::graph::NodeDefaults;
let graph = graph
.with_node_defaults(NodeDefaults::new().with_retry(RetryPolicy::new(3)))
.with_node_retry("critical", RetryPolicy::new(10));
Um valor por nΓ³ sempre prevalece. NodeDefaults tambΓ©m inclui um tempo limite e um manipulador de falhas.
Recuperando-se de uma falha de nΓ³
Quando o orΓ§amento de novas tentativas de um nΓ³ se esgota, um manipulador pode registrar o que aconteceu e nomear um nΓ³ de recuperaΓ§Γ£o em vez de encerrar a execuΓ§Γ£o:
let graph = graph.with_node_error_handler("charge", |node, error, _state| {
Ok(NodeOutput::new()
.with_update("status", json!(format!("{node} failed: {error}")))
.with_goto(["compensate"]))
});
Retornar Err encerra a execuΓ§Γ£o como antes. Uma interrupΓ§Γ£o nunca chega a um manipulador,
porque uma pausa nΓ£o Γ© uma falha.
Limitando o crescimento dos checkpoints
Uma thread acumula um checkpoint por superetapa. Portanto, uma execuΓ§Γ£o que permanece ativa por dias
cresce indefinidamente, o que aumenta o custo de armazenamento e torna list mais lento.
use adk_graph::checkpoint::RetentionPolicy;
use std::time::Duration;
let graph = graph
.with_checkpoint_retention(
RetentionPolicy::keep_last(50).with_max_age(Duration::from_secs(7 * 24 * 3600)),
);
A remoΓ§Γ£o ocorre apΓ³s cada salvamento, portanto o custo permanece proporcional Γ execuΓ§Γ£o e nenhum
job externo Γ© necessΓ‘rio. O checkpoint mais recente nunca Γ© descartado, independentemente da
polΓtica, porque Γ© aquele que um reinΓcio carrega β keep_last(0) Γ© elevado para
um, e uma thread cujo todos os checkpoints ultrapassaram o limite de idade mantΓ©m um.
Desativado por padrΓ£o, para que uma thread existente mantenha todo o seu histΓ³rico e a viagem no tempo ainda possa alcanΓ§ar cada etapa. Defina uma polΓtica quando uma thread tiver longa duraΓ§Γ£o e vocΓͺ nΓ£o precisar retroceder muito.
Rejeitando canais nΓ£o declarados
Um canal que o esquema nΓ£o declara usa o redutor de substituiΓ§Γ£o, porque esse Γ© o fallback para um nome desconhecido. Um grafo que declarou um canal de lista e depois escreveu um nome quase correspondente mantΓ©m apenas o ΓΊltimo valor e nΓ£o relata nada.
let graph = graph.with_strict_channels();
Um nΓ³ que escreve em um canal nΓ£o declarado falha a execuΓ§Γ£o com
GraphError::UndeclaredChannel, identificando o nΓ³ e o canal. Desativado por padrΓ£o
e inerte quando um grafo nΓ£o declara nenhum canal.
IntegraΓ§Γ£o com ADK
GraphAgent implementa o trait Agent de ADK, portanto funciona com:
- Executor: Use com
adk-runnerpara execuΓ§Γ£o padrΓ£o - Callbacks: Suporte completo para callbacks antes/depois
- SessΓ΅es: Funciona com
adk-sessionpara o histΓ³rico da conversa - Streaming: Retorna ADK
EventStream
use adk_runner::Runner;
let graph_agent = GraphAgent::builder("workflow")
.before_agent_callback(|ctx| async {
println!("Starting graph execution for session: {}", ctx.session_id());
Ok(())
})
.after_agent_callback(|ctx, event| async {
if let Some(content) = event.content() {
println!("Graph completed with content");
}
Ok(())
})
// ... graph definition
.build()?;
// GraphAgent implements Agent trait - use with Launcher or Runner
// See adk-runner README for Runner configuration
Exemplos
Exemplos de grafos validados neste repositΓ³rio:
cargo run --manifest-path examples/tier_examples/standard/Cargo.toml --bin 11-standard-graph
cargo run --manifest-path examples/tier_examples/standard/Cargo.toml --bin 12-standard-sequential
cargo run --manifest-path examples/competitive_graph_resume/Cargo.toml
A galeria completa de grafos com integraΓ§Γ£o real de LLM estΓ‘ disponΓvel em adk-playground.
ComparaΓ§Γ£o com LangGraph
| Recurso | LangGraph | adk-graph |
|---|---|---|
| GestΓ£o de estado | TypedDict + redutores | StateSchema + redutores |
| Modelo de execuΓ§Γ£o | superetapas do Pregel | superetapas do Pregel |
| PersistΓͺncia de checkpoints | MemΓ³ria, SQLite, Postgres | MemΓ³ria, SQLite |
| Humano no circuito | interrupt_before/after | interrupt_before/after + dinΓ’mico |
| Streaming | 5 modos | 5 modos |
| Ciclos | Nativo | Nativo |
| SeguranΓ§a de tipos | Tipagem do Python | Sistema de tipos do Rust |
| IntegraΓ§Γ£o com LLM | LangChain | Agentes AgentNode + ADK |
| Roteamento a partir de um nΓ³ | Command(goto=...) | NodeOutput::with_goto, AgentNode::with_goto_mapper |
| ExpansΓ£o dimensionada pelo estado | Send("node", input) | ctx.run_node_with(name, input, options) |
| Nova tentativa por nΓ³ | RetryPolicy | RetryPolicy com backoff limitado e jitter |
| Limite de concorrΓͺncia | max_concurrency na configuraΓ§Γ£o | with_max_concurrency no grafo |
| Cache de nΓ³ | cache_policy | cache_policy (recurso node-cache) |
| JunΓ§Γ£o adiada | defer=True | AutomΓ‘tica para vΓ‘rios predecessores diretos, alΓ©m de um quΓ³rum de n entre m |
| Subgrafo como um nΓ³ | add_node("sub", compiled) | SubgraphNode, com o mapeamento de canais verificado quando o pai Γ© compilado |
| Salto do subgrafo para o pai | Command(graph=Command.PARENT) | NodeOutput::with_goto_parent |
| PadrΓ΅es de nΓ³s em todo o grafo | set_node_defaults (β₯1.2) | with_node_defaults, alΓ©m do default_timeout preexistente |
| Manipuladores de falha de nΓ³s | error_handler (β₯1.2) | with_node_error_handler, executado quando o orΓ§amento de novas tentativas se esgota |
Vale destacar claramente duas diferenΓ§as. run_node_with retorna o resultado do filho em linha e o registra, portanto uma execuΓ§Γ£o retomada nΓ£o paga novamente por um filho concluΓdo, enquanto Send encaminha o trabalho ao agendador e o coleta por meio de um redutor. AlΓ©m disso, ainda nΓ£o hΓ‘ aqui um equivalente para um salto de um subgrafo para seu pai.
Anterior: β Sistemas multiagente | PrΓ³ximo: Agentes em tempo real β