Agentes de Grafo
Crie fluxos de trabalho complexos e com estado usando orquestraรงรฃo no estilo LangGraph com integraรงรฃo nativa ADK-Rust.
Visรฃo geral
GraphAgent permite definir fluxos de trabalho como grafos direcionados com nรณs e arestas, oferecendo suporte a:
- AgentNode: Envolva agentes LLM como nรณs do grafo com mapeadores de entrada/saรญda personalizados
- Fluxos de Trabalho Cรญclicos: Suporte nativo para 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 reducers (sobrescrever, anexar, somar, personalizado)
- Checkpointing: Estado persistente para tolerรขncia a falhas e humano-no-loop
- Streaming: Vรกrios modos de stream (values, updates, messages, debug)
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 fluxo de trabalho baseados em grafo para o ecossistema ADK-Rust, mantendo compatibilidade total com o sistema de agentes de ADK.
Principais Benefรญcios:
- Design Visual de Fluxo de Trabalho: Defina lรณgica complexa como grafos intuitivos de nรณs e arestas
- Execuรงรฃo Paralela: Vรกrios nรณs podem ser executados simultaneamente para melhor desempenho
- Persistรชncia de Estado: Checkpointing integrado para tolerรขncia a falhas e humano-no-loop
- Integraรงรฃo LLM: Suporte nativo para envolver agentes ADK como nรณs do grafo
- Roteamento Flexรญvel: Arestas estรกticas, roteamento condicional e tomada de decisรฃo dinรขmica
O que vocรช vai construir
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 executam trabalho (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 nรณs
- Execuรงรฃo Paralela - Vรกrios nรณs podem ser executados simultaneamente para melhor desempenho
Entendendo os Componentes Principais
๐ง Nรณs: Os Trabalhadores Os nรณs sรฃo onde o trabalho real acontece. Cada nรณ pode:
- AgentNode: Envolver 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 funรงรฃo e expertise especรญficas.
๐ Arestas: O Controle de Fluxo As arestas determinam como a execuรงรฃo avanรงa 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 de estrada que direcionam o fluxo de trabalho.
๐พ Estado: A Memรณria Compartilhada O estado รฉ um armazenamento de chave-valor que todos os nรณs podem ler e escrever:
- Dados de Entrada: Informaรงรตes iniciais alimentadas no 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 age como um quadro branco compartilhado onde os nรณs podem deixar informaรงรตes para outros usarem.
โก Execuรงรฃo Paralela: O Aumento de Velocidade Quando vรกrias arestas saem de um nรณ, esses 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 I/O
- Escalabilidade: Lide com fluxos de trabalho mais complexos sem desaceleraรงรฃo linear
Isso รฉ como ter vรกrios trabalhadores enfrentando diferentes partes de um trabalho ao mesmo tempo, em vez de esperar na fila.
Inรญcio Rรกpido
1. Crie Seu Projeto
cargo new graph_demo
cd graph_demo
Adicione dependรชncias em Cargo.toml:
[dependencies]
adk-graph = { version = "2.0.0", features = ["sqlite"] }
adk-agent = "2.0.0"
adk-model = "2.0.0"
adk-core = "2.0.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
Aqui estรก 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-2.5-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 a Execuรงรฃo do Grafo Funciona
A Visรฃo Geral
Agentes de grafo executam em superpassos - todos os nรณs prontos sรฃo executados em paralelo, depois o grafo aguarda a conclusรฃo de todos antes do prรณximo passo:
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 de Estado Atravรฉs dos Nรณs
Cada nรณ pode ler e escrever 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 | Encapsula agentes LLM com mapeadores de entrada/saรญda |
input_mapper | Transforma 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 nรณs |
ExecutionConfig | Fornece o ID da thread para checkpointing |
Roteamento Condicional com Classificaรงรฃo LLM
Construa sistemas de roteamento inteligentes em que LLMs decidem o caminho de execuรงรฃo:
Visual: Roteamento Baseado em Sentimento
โโโโโโโโโโโโโโโโโโโโโโโ
User Feedback โ โ
โโโโโโโโโโโโโโโโโถ โ CLASSIFIER โ
โ ๐ง Analyze tone โ
โโโโโโโโโโโโฌโโโโโโโโโโโ
โ
โโโโโโโโโโโโโโโโโผโโโโโโโโโโโโโโโโ
โ โ โ
โผ โผ โผ
โโโโโโโโโโโโโโโโโโโโ โโโโโโโโโโโโโโโโโโโโ โโโโโโโโโโโโโโโโโโโโ
โ POSITIVE โ โ NEGATIVE โ โ NEUTRAL โ
โ โ โ โ โ โ
โ ๐ Thank you! โ โ ๐ Apologize โ โ ๐ Ask more โ
โ Celebrate โ โ Help fix โ โ questions โ
โโโโโโโโโโโโโโโโโโโโ โโโโโโโโโโโโโโโโโโโโ โโโโโโโโโโโโโโโโโโโโ
Cรณdigo Completo de 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-2.5-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 de 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
Construa agentes que podem usar ferramentas iterativamente para resolver problemas complexos:
Visual: Ciclo de 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-2.5-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 de 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 (tipicamente LlmAgent) como um nรณ do grafo:
O que o agente vรช
Um agente dentro de um grafo รฉ executado sob um contexto derivado da invocaรงรฃo que iniciou o grafo, entรฃo ele se comporta da mesma forma que fora dele. Quando uma Runner invoca uma GraphAgent, a identidade e os serviรงos do chamador sรฃo transmitidos automaticamente:
| Transmitido | Observaรงรฃo |
|---|---|
app_name, user_id, session_id | Do chamador, nรฃo um sintรฉtico |
| Escopos e metadados da solicitaรงรฃo | Assim, as verificaรงรตes de escopo veem as permissรตes do chamador |
| Serviรงo secreto, memรณria, artefatos, estado compartilhado | Disponรญvel 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}, entรฃo os eventos de um nรณ sรฃo atribuรญveis |
Um grafo invocado diretamente โ graph.invoke(state, ExecutionConfig::new("thread")) โ
nรฃo tem uma invocaรงรฃo da qual herdar. Isso รฉ o modo standalone: o nรณ recebe
user_id = "graph_user", app_name = "graph_app", a branch main, sem segredos e sem
memรณria. ร um modo deliberado para executar um grafo fora de um Runner, nรฃo uma
alternativa de fallback para usar em produรงรฃo.
Para fazer a ponte manualmente โ por exemplo, ao acionar 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());
Nota: o nรณ ainda รฉ executado por conta prรณpria em uma sessรฃo de grafo na memรณria, entรฃo o histรณrico de conversa do agente dentro de um nรณ fica restrito ao nรณ, em vez de ser anexado ร 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),
],
)
Auxiliares de Roteador
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()?;
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 especializados:
// 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 Reducers
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 Reducer
| Redutor | Comportamento |
|---|---|
Overwrite | Substitui o valor antigo pelo novo (padrรฃo) |
Append | Adiciona ร lista |
Sum | Adicionar valores numรฉricos |
Custom | Funรงรฃo de mesclagem personalizada |
Checkpointing
Ative o estado persistente para tolerรขncia a falhas e intervenรงรฃo humana no loop:
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 checkpoint registra
Um checkpoint armazena o estado acumulado, o nรบmero da etapa e a frontier โ os nรณs que ainda precisam ser executados. Ele รฉ gravado depois que a frontier avanรงa, entรฃo retomar nunca reexecuta um nรณ que jรก foi concluรญdo e nunca aplica suas atualizaรงรตes duas vezes. Uma execuรงรฃo que termina faz checkpoint de uma frontier vazia, entรฃo retomar um thread concluรญdo retorna o estado final em vez de reiniciar o grafo.
Dois casos fazem checkpoint deliberadamente da frontier 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 na retomada:
| Situaรงรฃo | Fronteira salva |
|---|---|
| Super-etapa concluรญda | Os prรณximos nรณs a executar |
| Execuรงรฃo concluรญda | Vazio |
| Interrupรงรฃo acionada (bloqueando ou em streaming) | Os nรณs que estavam em execuรงรฃo |
Os checkpoints de execuรงรตes em streaming seguem o mesmo cronograma das execuรงรตes bloqueantes, inclusive quando uma interrupรงรฃo encerra o stream, entรฃo uma pausa com humano no loop pode ser retomada em qualquer modo de execuรงรฃo.
Histรณrico de Checkpoints (Viagem no Tempo)
Somente leitura.
TimeTravelHandle::state_history(from, to)retorna o estado que foi armazenado em cada etapa com checkpoint. 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 no histรณrico, usefork_atpara criar um ramo desse checkpoint e invocar o grafo na thread bifurcada. O mรฉtodo antes se chamavareplaye era documentado como reexecuรงรฃo do grafo, o que ele nunca fez.
Os checkpoints tambรฉm habilitam retomada durรกvel โ se uma execuรงรฃo do grafo falhar ou o processo reiniciar, a execuรงรฃo continua a partir do รบltimo checkpoint persistido em vez de recomeรงar do inรญcio. Use SqliteCheckpointer ou PostgresCheckpointer para persistรชncia segura contra falhas.
// 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);
}
Humano no Loop
Pause a execuรงรฃo para 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
Envie eventos em stream conforme 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 Stream
| 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 |
Messages lรช tokens de Node::execute_stream ร medida que sรฃo produzidos.
Cada nรณ รฉ executado uma vez por super-step nesse modo: o nรณ relata suas atualizaรงรตes de estado no stream 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 faturรกvel ao modelo por nรณ.
Importante: um
Nodepersonalizado que sobrescreveexecute_streamdeve emitir 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 para vocรช.
As polรญticas de timeout se aplicam ร prรณpria execuรงรฃo em stream. Para um stream,
idle_timeout significa que nenhum evento foi produzido dentro do limite.
ADK Integraรงรฃo
GraphAgent implementa o trait ADK Agent, entรฃo funciona com:
- Runner: Use com
adk-runnerpara execuรงรฃo padrรฃo - Callbacks: Suporte completo para callbacks before/after
- Sessions: Funciona com
adk-sessionpara histรณrico de 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 graph 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
O ADK-Rust Playground incorporado a este site inclui a galeria completa de grafos com integraรงรฃo real de LLM.
Comparaรงรฃo com LangGraph
| Recurso | LangGraph | adk-graph |
|---|---|---|
| Gerenciamento de estado | TypedDict + Redutores | StateSchema + Redutores |
| Modelo de execuรงรฃo | Pregel super-steps | Pregel super-steps |
| Pontos de controle | Memรณria, SQLite, Postgres | Memรณria, SQLite |
| Humano no loop | interrupt_before/after | interrupt_before/after + dinรขmico |
| Streaming | 5 modos | 5 modos |
| Ciclos | Suporte nativo | Suporte nativo |
| Seguranรงa de tipos | tipagem Python | sistema de tipos Rust |
| LLM Integraรงรฃo | LangChain | AgentNode + ADK agentes |
Anterior: โ Sistemas Multiagente | Prรณximo: Agentes em Tempo Real โ