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

ComponenteFunรงรฃo
AgentNodeEncapsula agentes LLM com mapeadores de entrada/saรญda
input_mapperTransforma estado โ†’ entrada do agente Content
output_mapperTransforma eventos do agente โ†’ atualizaรงรตes de estado
channelsDeclara os campos de estado que o grafo usarรก
edge()Define o fluxo de execuรงรฃo entre nรณs
ExecutionConfigFornece 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:

TransmitidoObservaรงรฃo
app_name, user_id, session_idDo chamador, nรฃo um sintรฉtico
Escopos e metadados da solicitaรงรฃoAssim, as verificaรงรตes de escopo veem as permissรตes do chamador
Serviรงo secreto, memรณria, artefatos, estado compartilhadoDisponรญvel exatamente como fora do grafo
CancelamentoRunner::interrupt alcanรงa um agente em execuรงรฃo como um nรณ
RunConfigHerdado do chamador
branchDerivado, 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

RedutorComportamento
OverwriteSubstitui o valor antigo pelo novo (padrรฃo)
AppendAdiciona ร  lista
SumAdicionar valores numรฉricos
CustomFunรงรฃ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รงรฃoFronteira salva
Super-etapa concluรญdaOs prรณximos nรณs a executar
Execuรงรฃo concluรญdaVazio
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, use fork_at para criar um ramo desse checkpoint e invocar o grafo na thread bifurcada. O mรฉtodo antes se chamava replay e 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

ModoDescriรงรฃo
ValuesTransmite o estado completo apรณs cada nรณ
UpdatesTransmite apenas as alteraรงรตes de estado
MessagesTransmitir atualizaรงรตes de tipo de mensagem
DebugTransmitir 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 Node personalizado que sobrescreve execute_stream deve emitir um evento StreamEvent::Updates contendo suas atualizaรงรตes de estado. Sem isso, o nรณ transmite eventos, mas nรฃo contribui com estado no modo Messages. O execute_stream padrรฃo, que encapsula execute, 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-runner para execuรงรฃo padrรฃo
  • Callbacks: Suporte completo para callbacks before/after
  • Sessions: Funciona com adk-session para 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

RecursoLangGraphadk-graph
Gerenciamento de estadoTypedDict + RedutoresStateSchema + Redutores
Modelo de execuรงรฃoPregel super-stepsPregel super-steps
Pontos de controleMemรณria, SQLite, PostgresMemรณria, SQLite
Humano no loopinterrupt_before/afterinterrupt_before/after + dinรขmico
Streaming5 modos5 modos
CiclosSuporte nativoSuporte nativo
Seguranรงa de tipostipagem Pythonsistema de tipos Rust
LLM IntegraรงรฃoLangChainAgentNode + ADK agentes

Anterior: โ† Sistemas Multiagente | Prรณximo: Agentes em Tempo Real โ†’

Agentes de Grafo - Documentaรงรฃo ADK-Rust | ADK-Rust