Agentes de grafo

Construye flujos de trabajo complejos y con estado usando orquestación estilo LangGraph con integración nativa de ADK-Rust.

Resumen

GraphAgent te permite definir flujos de trabajo como grafos dirigidos con nodos y aristas, con soporte para:

  • AgentNode: Envuelve agentes LLM como nodos del grafo con mapeadores de entrada/salida personalizados
  • Flujos de trabajo cíclicos: Soporte nativo para bucles y razonamiento iterativo (patrón ReAct)
  • Enrutamiento condicional: Enrutamiento dinámico de aristas basado en el estado
  • Gestión de estado: Estado tipado con reducers (sobrescribir, añadir, sumar, personalizado)
  • Checkpointing: Estado persistente para tolerancia a fallos y participación humana en el proceso
  • Streaming: Múltiples modos de flujo (values, updates, messages, debug)

La crate adk-graph proporciona orquestación de flujos de trabajo estilo LangGraph para construir flujos de trabajo de agentes complejos y con estado. Aporta capacidades de flujo de trabajo basadas en grafos al ecosistema ADK-Rust manteniendo compatibilidad total con el sistema de agentes de ADK.

Beneficios clave:

  • Diseño visual del flujo de trabajo: Define lógica compleja como grafos intuitivos de nodos y aristas
  • Ejecución paralela: Varios nodos pueden ejecutarse simultáneamente para mejor rendimiento
  • Persistencia de estado: Checkpointing integrado para tolerancia a fallos y participación humana en el proceso
  • Integración con LLM: Soporte nativo para envolver agentes ADK como nodos del grafo
  • Enrutamiento flexible: Aristas estáticas, enrutamiento condicional y toma de decisiones dinámica

Lo que construirás

En esta guía, crearás una tubería de procesamiento de texto que ejecuta traducción y resumen en paralelo:

                        ┌─────────────────────┐
       User Input       │                     │
      ────────────────▶ │       START         │
                        │                     │
                        └──────────┬──────────┘
                                   │
                   ┌───────────────┴───────────────┐
                   │                               │
                   ▼                               ▼
        ┌──────────────────┐            ┌──────────────────┐
        │   TRANSLATOR     │            │   SUMMARIZER     │
        │                  │            │                  │
        │  🇫🇷 French       │            │  📝 One sentence │
        │     Translation  │            │     Summary      │
        └─────────┬────────┘            └─────────┬────────┘
                  │                               │
                  └───────────────┬───────────────┘
                                  │
                                  ▼
                        ┌─────────────────────┐
                        │      COMBINE        │
                        │                     │
                        │  📋 Merge Results   │
                        └──────────┬──────────┘
                                   │
                                   ▼
                        ┌─────────────────────┐
                        │        END          │
                        │                     │
                        │   ✅ Complete       │
                        └─────────────────────┘

Conceptos clave:

  • Nodos - Unidades de procesamiento que realizan trabajo (agentes LLM, funciones o lógica personalizada)
  • Aristas - Flujo de control entre nodos (conexiones estáticas o enrutamiento condicional)
  • Estado - Datos compartidos que fluyen por el grafo y persisten entre nodos
  • Ejecución paralela - Varios nodos pueden ejecutarse simultáneamente para mejor rendimiento

Entender los componentes principales

🔧 Nodos: Los trabajadores Los nodos son donde ocurre el trabajo real. Cada nodo puede:

  • AgentNode: Envolver un agente LLM para procesar lenguaje natural
  • Nodo de función: Ejecutar código Rust personalizado para el procesamiento de datos
  • Nodos integrados: Usar lógica predefinida como contadores o validadores

Piensa en los nodos como trabajadores especializados en una línea de ensamblaje: cada uno tiene una tarea y experiencia específicas.

🔀 Aristas: El control del flujo Las aristas determinan cómo avanza la ejecución a través de tu grafo:

  • Aristas estáticas: Conexiones directas (A → B → C)
  • Aristas condicionales: Enrutamiento dinámico basado en el estado (if sentiment == "positive" → positive_handler)
  • Aristas paralelas: Múltiples rutas desde un nodo (START → [translator, summarizer])

Las aristas son como semáforos y señales de tránsito que dirigen el flujo del trabajo.

💾 Estado: La memoria compartida El estado es un almacén clave-valor desde el que todos los nodos pueden leer y escribir:

  • Datos de entrada: Información inicial introducida en el grafo
  • Resultados intermedios: La salida de un nodo se convierte en la entrada de otro
  • Salida final: El resultado completado después de todo el procesamiento

El estado actúa como una pizarra compartida donde los nodos pueden dejar información para que otros la usen.

⚡ Ejecución paralela: El impulso de velocidad Cuando varias aristas salen de un nodo, esos nodos destino se ejecutan simultáneamente:

  • Procesamiento más rápido: Las tareas independientes se ejecutan al mismo tiempo
  • Eficiencia de recursos: Mejor aprovechamiento de CPU y E/S
  • Escalabilidad: Maneja flujos de trabajo más complejos sin desaceleración lineal

Esto es como tener varios trabajadores abordando distintas partes de un trabajo al mismo tiempo en lugar de esperar en fila.


Inicio rápido

1. Crea tu proyecto

cargo new graph_demo
cd graph_demo

Añade dependencias a 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"

Crea .env con tu clave de API:

echo 'GOOGLE_API_KEY=your-api-key' > .env

2. Ejemplo de procesamiento paralelo

Aquí tienes un ejemplo completo y funcional que procesa texto en 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(())
}

Salida del ejemplo:

=== 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.

Cómo funciona la ejecución del grafo

La visión general

Los agentes de grafo se ejecutan en superpasos: todos los nodos listos se ejecutan en paralelo, luego el grafo espera a que todos terminen antes del siguiente paso:

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 ✅

Flujo del estado a través de los nodos

Cada nodo puede leer y escribir en el estado compartido:

┌─────────────────────────────────────────────────────────────────────┐
│ 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..."          │
│   }                                                                 │
│                                                                     │
└─────────────────────────────────────────────────────────────────────┘

Qué lo hace funcionar

ComponenteFunción
AgentNodeEncapsula agentes LLM con mapeadores de entrada/salida
input_mapperTransforma el estado → entrada del agente Content
output_mapperTransforma los eventos del agente → actualizaciones de estado
channelsDeclara los campos de estado que utilizará el grafo
edge()Define el flujo de ejecución entre nodos
ExecutionConfigProporciona el ID del hilo para el punto de control

Enrutamiento condicional con clasificación LLM

Construye sistemas de enrutamiento inteligentes donde LLMs deciden la ruta de ejecución:

Visual: Enrutamiento basado en sentimiento

                        ┌─────────────────────┐
       User Feedback    │                     │
      ────────────────▶ │    CLASSIFIER       │
                        │  🧠 Analyze tone    │
                        └──────────┬──────────┘
                                   │
                   ┌───────────────┼───────────────┐
                   │               │               │
                   ▼               ▼               ▼
        ┌──────────────────┐ ┌──────────────────┐ ┌──────────────────┐
        │   POSITIVE       │ │    NEGATIVE      │ │    NEUTRAL       │
        │                  │ │                  │ │                  │
        │  😊 Thank you!   │ │  😔 Apologize    │ │  😐 Ask more     │
        │     Celebrate    │ │     Help fix     │ │     questions    │
        └──────────────────┘ └──────────────────┘ └──────────────────┘

Código de ejemplo completo

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(())
}

Flujo de ejemplo:

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?"

Patrón ReAct: razonamiento + acción

Construye agentes que puedan usar herramientas iterativamente para resolver problemas complejos:

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)  │
                        └─────────────────────┘

Ejemplo 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(())
}

Flujo de ejemplo:

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

Envuelve cualquier ADK Agent (normalmente LlmAgent) como un nodo de grafo:

Lo que ve el agente

Un agente dentro de un grafo se ejecuta bajo un contexto derivado de la invocación que inició el grafo, por lo que se comporta igual que fuera de uno. Cuando una Runner invoca una GraphAgent, la identidad y los servicios del llamador se transmiten automáticamente:

TransmitidoNota
app_name, user_id, session_idLa del llamador, no una sintética
Ámbitos y metadatos de la solicitudPara que las comprobaciones de ámbito vean las concesiones del llamador
Servicio secreto, memoria, artefactos, estado compartidoDisponible exactamente igual que fuera del grafo
CancelaciónRunner::interrupt llega a un agente que se ejecuta como un nodo
RunConfigHeredado del llamador
branchDerivado, como {caller_branch}.{agent_name}, para que los eventos de un nodo sean atribuibles

Un grafo invocado directamente — graph.invoke(state, ExecutionConfig::new("thread")) — no tiene ninguna invocación de la cual heredar. Eso es modo independiente: el nodo obtiene user_id = "graph_user", app_name = "graph_app", la rama main, sin secretos y sin memoria. Es un modo deliberado para ejecutar un grafo fuera de un Runner, no una alternativa de respaldo a la que recurrir en producción.

Para enlazarlo manualmente — por ejemplo, cuando se controla un grafo desde tu propio executor — pasa la invocación explícitamente:

let config = ExecutionConfig::new(ctx.session_id()).with_parent_context(ctx.clone());

Nota: el nodo sigue ejecutándose por su cuenta en una sesión de grafo en memoria, por lo que el historial de conversación del agente dentro de un nodo queda limitado al nodo en lugar de añadirse a la sesión del llamador.

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
    });

Nodos de función

Funciones asíncronas simples que procesan el 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 aristas

Aristas estáticas

Conexiones directas entre nodos:

.edge(START, "first_node")
.edge("first_node", "second_node")
.edge("second_node", END)

Aristas condicionales

Enrutamiento dinámico basado en el 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),
    ],
)

Ayudantes de enrutamiento

Usa routers integrados para patrones comunes:

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),
])

Ejecución en paralelo

Varias aristas desde un solo nodo se ejecutan en 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 (patrón ReAct)

Construye agentes de razonamiento iterativo con 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

Dirige tareas a 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()?;

Gestión del estado

Esquema de estado con reductores

Controla cómo se combinan las actualizaciones del estado:

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 reductor

ReductorComportamiento
OverwriteReemplazar el valor antiguo por el nuevo (predeterminado)
AppendAñadir a la lista
SumAgregar valores numéricos
CustomFunción de fusión personalizada

Puntos de control

Habilita estado persistente para tolerancia a fallos y human-in-the-loop:

En memoria (desarrollo)

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 (producción)

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);

Qué registra un punto de control

Un punto de control almacena el estado acumulado, el número de paso y la frontera — los nodos que aún deben ejecutarse. Se escribe después de que la frontera avanza, por lo que reanudar nunca vuelve a ejecutar un nodo que ya se completó y nunca aplica dos veces sus actualizaciones. Una ejecución que termina guarda un punto de control con una frontera vacía, así que reanudar un hilo completado devuelve el estado final en lugar de reiniciar el grafo.

Dos casos guardan deliberadamente en el punto de control la frontera que se estaba ejecutando en lugar de la siguiente, porque el nodo interrumpido aún no ha producido sus actualizaciones y debe ejecutarse de nuevo al reanudar:

SituaciónFrontier guardada
Super-step completadoLos siguientes nodos a ejecutar
Ejecución finalizadaVacío
Interrupción generada (bloqueante o en streaming)Los nodos que se estaban ejecutando

Las ejecuciones transmitidas crean puntos de control con la misma programación que las ejecuciones bloqueantes, incluso cuando una interrupción termina la transmisión, de modo que una pausa con participación humana puede reanudarse en cualquiera de los dos modos de ejecución.

Historial de puntos de control (viaje en el tiempo)

Solo lectura. TimeTravelHandle::state_history(from, to) devuelve el estado que se almacenó en cada paso con punto de control. No ejecuta nada: no se ejecuta ningún nodo, no se regenera ningún evento y no se repite ningún efecto secundario. Para volver a ejecutar desde un punto en el historial, usa fork_at para ramificar ese punto de control e invocar el grafo en el hilo bifurcado. El método anteriormente se llamaba replay y se documentaba como si reejecutara el grafo, cosa que nunca hizo.

Los puntos de control también permiten la reanudación duradera: si una ejecución del grafo falla o el proceso se reinicia, la ejecución se reanuda desde el último punto de control persistido en lugar de comenzar de nuevo. Usa SqliteCheckpointer o PostgresCheckpointer para una persistencia segura ante fallos.

// 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);
}

Participación humana en el ciclo

Pausa la ejecución para la aprobación humana usando interrupciones 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);
    }
}

Interrupciones estáticas

Usa interrupt_before o interrupt_after para puntos de pausa obligatorios:

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

Ejecución en streaming

Transmite eventos a medida que el grafo se ejecuta:

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 transmisión

ModoDescripción
ValuesTransmitir el estado completo después de cada nodo
UpdatesTransmitir solo los cambios de estado
MessagesTransmitir actualizaciones de tipo de mensaje
DebugTransmitir todos los eventos internos

Messages modo lee tokens de Node::execute_stream a medida que se producen. Cada nodo se ejecuta una vez por superpaso en este modo: el nodo informa sus actualizaciones de estado en el stream como un evento StreamEvent::Updates, y el ejecutor aplica esas actualizaciones en lugar de ejecutar el nodo una segunda vez para recopilarlas. Esto es especialmente importante para AgentNode, donde una segunda ejecución implicaría una segunda llamada facturada al modelo por nodo.

Importante: un Node personalizado que sobrescriba execute_stream debe emitir un evento StreamEvent::Updates que lleve sus actualizaciones de estado. Sin ello, el nodo transmite eventos pero no aporta estado en modo Messages. El execute_stream predeterminado, que envuelve execute, hace esto por ti.

Las políticas de tiempo de espera se aplican a la propia ejecución transmitida. Para un stream, idle_timeout significa que no se produjo ningún evento dentro del límite.

Integración de ADK

GraphAgent implementa el trait ADK Agent, por lo que funciona con:

  • Ejecutor: Úsalo con adk-runner para la ejecución estándar
  • Callbacks: Compatibilidad completa con callbacks antes/después
  • Sesiones: Funciona con adk-session para el historial de conversación
  • Streaming: Devuelve 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

Ejemplos

Ejemplos de grafos validados en este repositorio:

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

El ADK-Rust Playground integrado en este sitio web incluye la galería completa de grafos con integración real de LLM.

Comparación con LangGraph

FunciónLangGraphadk-graph
Gestión del estadoTypedDict + ReducersStateSchema + Reducers
Modelo de ejecuciónsuperpasos de Pregelsuperpasos de Pregel
Puntos de controlMemoria, SQLite, PostgresMemoria, SQLite
Humano en el ciclointerrupt_before/despuésinterrupt_before/después + dinámico
Streaming5 modos5 modos
CiclosSoporte nativoSoporte nativo
Seguridad de tiposTipado en PythonSistema de tipos de Rust
LLM IntegraciónLangChainAgentNode + ADK agentes

Anterior: ← Sistemas multiagente | Siguiente: Agentes en tiempo real →