Graph-Agenten

Erstellen Sie komplexe, zustandsbehaftete Workflows mit LangGraph-Γ€hnlicher Orchestrierung und nativer ADK-Rust-Integration.

Übersicht

GraphAgent ermΓΆglicht es Ihnen, Workflows als gerichtete Graphen mit Knoten und Kanten zu definieren und unterstΓΌtzt dabei:

  • AgentNode: LLM-Agenten als Graph-Knoten mit benutzerdefinierten Ein-/Ausgabe-Mappern einbinden
  • Zyklische Workflows: Native UnterstΓΌtzung fΓΌr Schleifen und iterative Schlussfolgerung (ReAct-Muster)
  • Bedingtes Routing: Dynamisches Kanten-Routing basierend auf dem Zustand
  • Zustandsverwaltung: Typisierter Zustand mit Reducern (ΓΌberschreiben, anhΓ€ngen, summieren, benutzerdefiniert)
  • Checkpointing: Persistenter Zustand fΓΌr Fehlertoleranz und Human-in-the-Loop
  • Streaming: Mehrere Stream-Modi (Werte, Aktualisierungen, Nachrichten, Debug)

Das adk-graph-Crate bietet LangGraph-artige Workflow-Orchestrierung zum Erstellen komplexer, zustandsbehafteter Agent-Workflows. Es bringt graphbasierte Workflow-Funktionen in das ADK-Rust-Γ–kosystem und bleibt dabei vollstΓ€ndig kompatibel mit dem Agentensystem von ADK.

Wichtige Vorteile:

  • Visuelles Workflow-Design: Definieren Sie komplexe Logik als intuitive Knoten-und-Kanten-Graphen
  • Parallele AusfΓΌhrung: Mehrere Knoten kΓΆnnen gleichzeitig ausgefΓΌhrt werden, fΓΌr bessere Leistung
  • Zustandspersistenz: Integriertes Checkpointing fΓΌr Fehlertoleranz und Human-in-the-Loop
  • LLM-Integration: Native UnterstΓΌtzung zum Einbinden von ADK-Agenten als Graph-Knoten
  • Flexibles Routing: Statische Kanten, bedingtes Routing und dynamische Entscheidungsfindung

Was Sie erstellen werden

In diesem Leitfaden erstellen Sie eine Textverarbeitungspipeline, die Übersetzung und Zusammenfassung parallel ausführt:

                        β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
       User Input       β”‚                     β”‚
      ────────────────▢ β”‚       START         β”‚
                        β”‚                     β”‚
                        β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
                                   β”‚
                   β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
                   β”‚                               β”‚
                   β–Ό                               β–Ό
        β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”            β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
        β”‚   TRANSLATOR     β”‚            β”‚   SUMMARIZER     β”‚
        β”‚                  β”‚            β”‚                  β”‚
        β”‚  πŸ‡«πŸ‡· French       β”‚            β”‚  πŸ“ One sentence β”‚
        β”‚     Translation  β”‚            β”‚     Summary      β”‚
        β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”˜            β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”˜
                  β”‚                               β”‚
                  β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
                                  β”‚
                                  β–Ό
                        β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
                        β”‚      COMBINE        β”‚
                        β”‚                     β”‚
                        β”‚  πŸ“‹ Merge Results   β”‚
                        β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
                                   β”‚
                                   β–Ό
                        β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
                        β”‚        END          β”‚
                        β”‚                     β”‚
                        β”‚   βœ… Complete       β”‚
                        β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜

Wichtige Konzepte:

  • Knoten - Verarbeitungseinheiten, die Arbeit ausfΓΌhren (LLM-Agenten, Funktionen oder benutzerdefinierte Logik)
  • Kanten - Kontrollfluss zwischen Knoten (statische Verbindungen oder bedingtes Routing)
  • Zustand - Gemeinsame Daten, die durch den Graphen fließen und zwischen Knoten erhalten bleiben
  • Parallele AusfΓΌhrung - Mehrere Knoten kΓΆnnen gleichzeitig ausgefΓΌhrt werden, fΓΌr bessere Leistung

Die Kernkomponenten verstehen

πŸ”§ Knoten: Die Arbeiter Knoten sind der Ort, an dem die eigentliche Arbeit stattfindet. Jeder Knoten kann:

  • AgentNode: Einen LLM-Agenten einbinden, um natΓΌrliche Sprache zu verarbeiten
  • Funktionsknoten: Benutzerdefinierten Rust-Code zur Datenverarbeitung ausfΓΌhren
  • Integrierte Knoten: Vorgegebene Logik wie ZΓ€hler oder Validierer verwenden

Stellen Sie sich Knoten als spezialisierte Arbeiter in einer Montagelinie vor - jeder hat eine bestimmte Aufgabe und Expertise.

πŸ”€ Kanten: Die Flusssteuerung Kanten bestimmen, wie die AusfΓΌhrung durch Ihren Graphen verlΓ€uft:

  • Statische Kanten: Direkte Verbindungen (A β†’ B β†’ C)
  • Bedingte Kanten: Dynamisches Routing basierend auf dem Zustand (if sentiment == "positive" β†’ positive_handler)
  • Parallele Kanten: Mehrere Pfade von einem Knoten (START β†’ [translator, summarizer])

Kanten sind wie Ampeln und Straßenschilder, die den Arbeitsfluss lenken.

πŸ’Ύ Zustand: Der gemeinsame Speicher Zustand ist ein Key-Value-Speicher, aus dem alle Knoten lesen und in den sie schreiben kΓΆnnen:

  • Eingabedaten: Anfangsinformationen, die in den Graphen eingespeist werden
  • Zwischenergebnisse: Die Ausgabe eines Knotens wird zur Eingabe fΓΌr einen anderen
  • Endausgabe: Das fertige Ergebnis nach der gesamten Verarbeitung

Zustand wirkt wie ein gemeinsames Whiteboard, auf dem Knoten Informationen fΓΌr andere hinterlassen kΓΆnnen.

⚑ Parallele Ausführung: Der Geschwindigkeitsschub Wenn mehrere Kanten einen Knoten verlassen, werden diese Zielknoten gleichzeitig ausgeführt:

  • Schnellere Verarbeitung: UnabhΓ€ngige Aufgaben laufen gleichzeitig
  • Ressourceneffizienz: Bessere Ausnutzung von CPU und I/O
  • Skalierbarkeit: BewΓ€ltigen Sie komplexere Workflows ohne lineare Verlangsamung

Das ist, als wΓΌrden mehrere Arbeiter verschiedene Teile eines Auftrags gleichzeitig bearbeiten, anstatt in einer Schlange zu warten.


Schnellstart

1. Erstellen Sie Ihr Projekt

cargo new graph_demo
cd graph_demo

FΓΌgen Sie AbhΓ€ngigkeiten zu Cargo.toml hinzu:

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

Erstellen Sie .env mit Ihrem API-SchlΓΌssel:

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

2. Beispiel fΓΌr parallele Verarbeitung

Hier ist ein vollstΓ€ndiges, lauffΓ€higes Beispiel, das Text parallel verarbeitet:

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

Beispielausgabe:

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

Wie die GraphausfΓΌhrung funktioniert

Das Gesamtbild

Graph-Agenten fΓΌhren AusfΓΌhrungen in Super-Schritten aus - alle bereiten Knoten laufen parallel, dann wartet der Graph, bis alle abgeschlossen sind, bevor der nΓ€chste Schritt beginnt:

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 βœ…

Datenfluss durch Knoten

Jeder Knoten kann aus dem gemeinsamen Zustand lesen und in ihn schreiben:

β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
β”‚ 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..."          β”‚
β”‚   }                                                                 β”‚
β”‚                                                                     β”‚
β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜

Was es mΓΆglich macht

KomponenteRolle
AgentNodeUmhΓΌllt LLM Agenten mit Ein-/Ausgabe-Mappern
input_mapperTransformiert Zustand β†’ Agenteneingabe Content
output_mapperTransformiert Agent-Ereignisse β†’ Zustandsaktualisierungen
channelsDeklariert Zustandsfelder, die der Graph verwenden wird
edge()Definiert den AusfΓΌhrungsfluss zwischen Knoten
ExecutionConfigStellt die Thread-ID fΓΌr das Checkpointing bereit

Bedingtes Routing mit LLM-Klassifizierung

Erstellen Sie intelligente Routing-Systeme, bei denen LLMs den AusfΓΌhrungspfad bestimmen:

Visualisierung: Routing basierend auf Stimmung

                        β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
       User Feedback    β”‚                     β”‚
      ────────────────▢ β”‚    CLASSIFIER       β”‚
                        β”‚  🧠 Analyze tone    β”‚
                        β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
                                   β”‚
                   β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”Όβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
                   β”‚               β”‚               β”‚
                   β–Ό               β–Ό               β–Ό
        β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β” β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
        β”‚   POSITIVE       β”‚ β”‚    NEGATIVE      β”‚ β”‚    NEUTRAL       β”‚
        β”‚                  β”‚ β”‚                  β”‚ β”‚                  β”‚
        β”‚  😊 Thank you!   β”‚ β”‚  πŸ˜” Apologize    β”‚ β”‚  😐 Ask more     β”‚
        β”‚     Celebrate    β”‚ β”‚     Help fix     β”‚ β”‚     questions    β”‚
        β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜

VollstΓ€ndiger Beispielcode

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

Beispielfluss:

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

ReAct-Muster: Denken + Handeln

Erstellen Sie Agenten, die Tools iterativ verwenden kΓΆnnen, um komplexe Probleme zu lΓΆsen:

Visualisierung: ReAct-Zyklus

                        β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
       User Question    β”‚                     β”‚
      ────────────────▢ β”‚      REASONER       β”‚
                        β”‚  🧠 Think + Act     β”‚
                        β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
                                   β”‚
                                   β–Ό
                        β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
                        β”‚   Has tool calls?   β”‚
                        β”‚                     β”‚
                        β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
                                   β”‚
                   β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
                   β”‚                               β”‚
                   β–Ό                               β–Ό
        β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”            β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
        β”‚       YES        β”‚            β”‚        NO        β”‚
        β”‚                  β”‚            β”‚                  β”‚
        β”‚  πŸ”„ Loop back    β”‚            β”‚  βœ… Final answer β”‚
        β”‚     to reasoner  β”‚            β”‚      END         β”‚
        β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”˜            β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
                  β”‚
                  └─────────────────┐
                                    β”‚
                                    β–Ό
                        β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
                        β”‚      REASONER       β”‚
                        β”‚  🧠 Think + Act     β”‚
                        β”‚   (next iteration)  β”‚
                        β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜

VollstΓ€ndiges ReAct-Beispiel

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

Beispielfluss:

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

Kapselt jeden ADK Agent (typischerweise LlmAgent) als Graphknoten:

Was der Agent sieht

Ein Agent innerhalb eines Graphen lÀuft unter einem Kontext, der von der Invocation abgeleitet ist, die den Graphen gestartet hat, sodass er sich genauso verhÀlt wie außerhalb eines solchen. Wenn ein Runner einen GraphAgent aufruft, werden die IdentitÀt und die Dienste des Aufrufers automatisch mitübertragen:

Übernommen durchHinweis
app_name, user_id, session_idDie des Aufrufers, nicht eine synthetische
Scopes und Request-MetadatenDamit Scope-PrΓΌfungen die Grants des Aufrufers sehen
Geheimdienst, Speicher, Artefakte, geteilter ZustandGenau wie außerhalb des Graphen verfügbar
AbbruchRunner::interrupt erreicht einen Agenten, der als Knoten ausgefΓΌhrt wird
RunConfigVom Aufrufer geerbt
branchAbgeleitet, als {caller_branch}.{agent_name}, sodass die Ereignisse eines Knotens zugeordnet werden kΓΆnnen

Ein Graph, der direkt aufgerufen wird β€” graph.invoke(state, ExecutionConfig::new("thread")) β€” hat keinen Aufruf, von dem er erben kann. Das ist der Standalone-Modus: Der Knoten erhΓ€lt user_id = "graph_user", app_name = "graph_app", Branch main, keine Geheimnisse und kein Speicher. Es ist ein bewusster Modus, um einen Graphen außerhalb von Runner auszufΓΌhren, nicht ein Fallback, auf den man in der Produktion zurΓΌckgreifen sollte.

Um die Verbindung manuell herzustellen β€” zum Beispiel, wenn Sie einen Graphen ΓΌber Ihren eigenen Executor steuern β€” ΓΌbergeben Sie den Aufruf explizit:

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

Hinweis: Der Knoten lΓ€uft weiterhin in seiner eigenen In-Memory-Graph-Session, sodass der GesprΓ€chsverlauf des Agenten innerhalb eines Knotens auf den Knoten beschrΓ€nkt ist und nicht an die Session des Aufrufers angehΓ€ngt wird.

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

Funktionsknoten

Einfache asynchrone Funktionen, die den Zustand verarbeiten:

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

Kantenarten

Statische Kanten

Direkte Verbindungen zwischen Knoten:

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

Bedingte Kanten

Dynamisches Routing basierend auf dem Zustand:

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

Router-Hilfen

Verwenden Sie integrierte Router fΓΌr gΓ€ngige Muster:

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

Parallele AusfΓΌhrung

Mehrere Kanten von einem einzelnen Knoten werden parallel ausgefΓΌhrt:

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

Zyklische Graphen (ReAct-Muster)

Erstellen Sie iterative Reasoning-Agenten mit Zyklen:

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

Multi-Agent-Supervisor

Leiten Sie Aufgaben an spezialisierte Agenten weiter:

// 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()?;

Zustandsverwaltung

Zustands-Schema mit Reducern

Steuern Sie, wie Zustandsaktualisierungen zusammengefΓΌhrt werden:

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

Reducer-Typen

ReducerVerhalten
OverwriteAlten Wert durch neuen ersetzen (Standard)
AppendAn Liste anhΓ€ngen
SumNumerische Werte hinzufΓΌgen
CustomBenutzerdefinierte ZusammenfΓΌhrungsfunktion

Checkpointing

Aktiviere persistenten Zustand fΓΌr Fehlertoleranz und Human-in-the-loop:

Im Speicher (Entwicklung)

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

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

Was ein Checkpoint erfasst

Ein Checkpoint speichert den akkumulierten Zustand, die Schrittzahl und die Frontier β€” die Knoten, die noch ausgefΓΌhrt werden mΓΌssen. Er wird geschrieben, nachdem sich die Frontier weiterbewegt hat, sodass das Fortsetzen nie einen Knoten erneut ausfΓΌhrt, der bereits abgeschlossen wurde, und seine Aktualisierungen nie doppelt anwendet. Ein Lauf, der abgeschlossen wird, checkpointet eine leere Frontier, sodass das Fortsetzen eines abgeschlossenen Threads den endgΓΌltigen Zustand zurΓΌckgibt, anstatt den Graphen neu zu starten.

Zwei FΓ€lle checkpointen absichtlich die Frontier, die gerade ausgefΓΌhrt wurde, statt der nΓ€chsten, weil der unterbrochene Knoten seine Aktualisierungen noch nicht erzeugt hat und beim Fortsetzen erneut ausgefΓΌhrt werden muss:

SituationFrontier gespeichert
Super-Schritt abgeschlossenDie nΓ€chsten auszufΓΌhrenden Knoten
AusfΓΌhrung abgeschlossenLeer
Interrupt ausgelΓΆst (blockierend oder streaming)Die Knoten, die ausgefΓΌhrt wurden

Streamed Runs checkpointen nach demselben Zeitplan wie blockierende Runs, auch wenn ein Interrupt den Stream beendet, sodass eine Pause mit Human-in-the-Loop in beiden AusfΓΌhrungsmodi fortgesetzt werden kann.

Checkpoint-Verlauf (Zeitreise)

Nur Lesezugriff. TimeTravelHandle::state_history(from, to) gibt den Zustand zurΓΌck, der an jedem checkpointierten Schritt gespeichert wurde. Es fΓΌhrt nichts aus β€” kein Node lΓ€uft, kein Event wird neu erzeugt, und kein Seiteneffekt wiederholt sich. Um ab einem Punkt in der Historie erneut auszufΓΌhren, verwende fork_at, um diesen Checkpoint zu verzweigen, und rufe den Graphen auf dem abgezweigten Thread auf. Die Methode hieß frΓΌher replay und wurde als erneute AusfΓΌhrung des Graphen dokumentiert, was sie nie getan hat.

Checkpoints ermΓΆglichen außerdem dauerhaftes Fortsetzen β€” wenn eine GraphausfΓΌhrung abstΓΌrzt oder der Prozess neu startet, wird die AusfΓΌhrung vom zuletzt persistent gespeicherten Checkpoint fortgesetzt, statt von vorn zu beginnen. Verwende SqliteCheckpointer oder PostgresCheckpointer fΓΌr ausfallsichere Persistenz.

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

Human-in-the-Loop

Pausiere die AusfΓΌhrung fΓΌr eine menschliche Freigabe mithilfe dynamischer Interrupts:

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

Statische Interrupts

Verwende interrupt_before oder interrupt_after fΓΌr obligatorische Pausenpunkte:

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

Streaming-AusfΓΌhrung

Streame Events, wΓ€hrend der Graph ausgefΓΌhrt wird:

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),
        _ => {}
    }
}

Stream-Modi

ModusBeschreibung
ValuesVollstΓ€ndigen Status nach jedem Knoten streamen
UpdatesNur StatusΓ€nderungen streamen
MessagesNachrichten-Typ-Aktualisierungen streamen
DebugAlle internen Ereignisse streamen

Messages-Modus liest Tokens aus Node::execute_stream, sobald sie erzeugt werden. Jeder Knoten lΓ€uft in diesem Modus pro Super-Schritt einmal: Der Knoten meldet seine Zustandsaktualisierungen im Stream als ein StreamEvent::Updates-Ereignis, und der Executor ΓΌbernimmt diese, statt den Knoten ein zweites Mal auszufΓΌhren, um sie zu sammeln. Das ist besonders wichtig fΓΌr AgentNode, wo eine zweite AusfΓΌhrung einen zweiten berechneten Model-Aufruf pro Knoten bedeuten wΓΌrde.

Wichtig: Ein benutzerdefinierter Node, der execute_stream überschreibt, muss ein StreamEvent::Updates-Ereignis liefern, das seine Zustandsaktualisierungen trÀgt. Ohne dieses streamt der Knoten Ereignisse, trÀgt aber im Messages-Modus keinen Zustand bei. Der standardmÀßige execute_stream, der execute umschließt, erledigt das für dich.

Timeout-Richtlinien gelten fΓΌr die gestreamte AusfΓΌhrung selbst. FΓΌr einen Stream bedeutet idle_timeout, dass innerhalb des Limits kein Ereignis erzeugt wurde.

ADK Integration

GraphAgent implementiert das ADK Agent-Trait und funktioniert daher mit:

  • Runner: Mit adk-runner fΓΌr die StandardausfΓΌhrung verwenden
  • Callbacks: VollstΓ€ndige UnterstΓΌtzung fΓΌr Callbacks vor/nach
  • Sessions: Funktioniert mit adk-session fΓΌr den GesprΓ€chsverlauf
  • Streaming: Gibt ADK EventStream zurΓΌck
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

Beispiele

Validierte Graph-Beispiele in diesem Repository:

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

Der in diese Website eingebettete ADK-Rust Playground enthΓ€lt die vollstΓ€ndige Graph-Galerie mit echter LLM-Integration.

Vergleich mit LangGraph

FunktionLangGraphadk-graph
ZustandsverwaltungTypedDict + ReducersStateSchema + Reducers
AusfΓΌhrungsmodellPregel-SuperstepsPregel-Supersteps
CheckpointingMemory, SQLite, PostgresMemory, SQLite
Mensch-im-Loopinterrupt_before/afterinterrupt_before/after + dynamisch
Streaming5 Modi5 Modi
ZyklenNative UnterstΓΌtzungNative UnterstΓΌtzung
TypsicherheitPython-TypisierungRust-Typsystem
LLM IntegrationLangChainAgentNode + ADK Agenten

Vorherige: ← Multi-Agent-Systeme | NΓ€chste: Echtzeit-Agenten β†’

Graph-Agenten - ADK-Rust Dokumentation | ADK-Rust