Funktionale API

Schreiben Sie Agenten-Workflows als normale asynchrone Rust-Funktionen mit automatischer Checkpoint-Erstellung, typisierten Zustandsreduzierern und Unterstützung für Unterbrechung und Fortsetzung.

Übersicht

Die funktionale API (die functional-Funktion in adk-graph) bietet eine übergeordnete Alternative zur expliziten Erstellung von Graphen, Knoten und Kanten. Anstatt eine StateGraph manuell zu erstellen, versehen Sie Funktionen mit den Makros #[entrypoint] und #[task] und verwenden den standardmäßigen Rust-Kontrollfluss.

Wichtige Vorteile:

  • Schreiben Sie Workflows als normale asynchrone Rust-Funktionen – keine Graph-DSL erforderlich
  • Automatische Checkpoint-Erstellung nach jeder Aufgabe zur Wiederherstellung nach Abstürzen
  • Standardmäßiger Rust-Kontrollfluss (if/else, for, loop, match)
  • Typisierte Zustandscontainer mit Garantien für die Persistenz
  • Kompatibel mit vorhandener StreamEvent- und Checkpointer-Infrastruktur

Erste Schritte

[dependencies]
adk-graph = { version = "2.1.0", features = ["functional"] }
adk-rust-macros = "2.1.0"

Zentrale Typen

TaskContext

Der Laufzeitkontext, der an alle Workflow-Funktionen übergeben wird. Bietet Zugriff auf Zustand, Checkpoint-Erstellung, Unterbrechungen und Streaming.

use adk_graph::functional::TaskContext;

// Read state
let count: Option<i64> = ctx.get("counter");

// Write state (uses configured reducer)
ctx.set("counter", serde_json::json!(count.unwrap_or(0) + 1));

// Emit progress events
ctx.emit(StreamEvent::custom("my_task", "progress", json!({"pct": 50})));

// Interrupt for human-in-the-loop
let approval: bool = ctx.interrupt("Please approve").await?;

ReducedValue<T>

Nur-anhängbarer Zustandscontainer, der über Checkpoints hinweg persistiert wird. Werte werden akkumuliert und niemals überschrieben.

use adk_graph::functional::ReducedValue;

let mut results: ReducedValue<String> = ReducedValue::new();
results.push("step 1 output".to_string());
results.push("step 2 output".to_string());
assert_eq!(results.len(), 2);
assert_eq!(&results[0], "step 1 output");

UntrackedValue<T>

Transiente Laufzeitwerte, die von der Checkpoint-Persistenz ausgeschlossen sind. Werden bei der Fortsetzung auf den Standardwert zurückgesetzt.

use adk_graph::functional::UntrackedValue;

let mut temp: UntrackedValue<Vec<u8>> = UntrackedValue::new();
temp.set(vec![1, 2, 3]);
// After checkpoint restore: temp.get() == &[]

MessagesValue

Container für Chatnachrichten mit automatischer Deduplizierung anhand von Nachrichten-IDs.

use adk_graph::functional::{MessagesValue, ChatMessage, MessageRole};

let mut messages = MessagesValue::new();
messages.push(ChatMessage {
    id: "msg-1".to_string(),
    role: MessageRole::User,
    content: "Hello".to_string(),
    metadata: None,
});
// Pushing same ID replaces the message (upsert)

StateSchemaValidator

Validiert Zustandstypen an den Grenzen des Workflows und erkennt Inkonsistenzen frühzeitig.

use adk_graph::functional::{StateSchemaValidator, ExpectedType};
use adk_graph::state::StateSchema;

let validator = StateSchemaValidator::new(schema)
    .expect_type("counter", ExpectedType::Number)
    .expect_type("status", ExpectedType::String)
    .require_field("status");

validator.validate_state(&state)?; // Fails with descriptive error

ExecutionLog

Verfolgt den Abschluss von Aufgaben für das Überspringen bei der Fortsetzung. Abgeschlossene Aufgaben werden bei der Fortsetzung eines Workflows übersprungen.

use adk_graph::functional::ExecutionLog;

let mut log = ExecutionLog::new();
log.record_start("fetch");
log.record_completion("fetch", json!({"data": [1,2,3]}));

// On resume: skip completed tasks
if log.is_completed("fetch") {
    let cached = log.get_result("fetch"); // Returns cached result
}

Hintergrundausführungen

Die background-Funktion in adk-server fügt REST-Endpunkte für die asynchrone Workflow-Ausführung hinzu.

Endpunkte

MethodePfadBeschreibung
POST/runsEinen Hintergrundlauf übermitteln
GET/runs/{run_id}Den Status des Laufs abfragen
LÖSCHEN/runs/{run_id}Ausführung abbrechen

Verwendung

Ein Lauf benennt einen workflowId. Etwas muss diesen Namen in Arbeit umwandeln; registrieren Sie daher einen WorkflowExecutor — ohne einen solchen schlägt ein übermittelter Lauf fehl, anstatt den Abschluss zu melden:

use adk_server::background::{
    BackgroundState, WorkflowRegistry, background_runs_router_with_state,
};
use serde_json::json;
use std::sync::Arc;

let registry = WorkflowRegistry::new().register("summarize", |input, cancel| async move {
    if cancel.is_cancelled() {
        return Err("cancelled before starting".to_string());
    }
    Ok(json!({ "summary": "…", "of": input.get("document") }))
});

let state = BackgroundState::new().with_executor(Arc::new(registry));
let app = axum::Router::new().merge(background_runs_router_with_state(state));

Das Übermitteln eines nicht registrierten workflowId gibt 404 zurück, anstatt einen Lauf einzureihen, der niemals ausgeführt werden kann. Implementieren Sie WorkflowExecutor direkt, um eine Verbindung zu adk-graph, dem funktionalen API oder Ihrem eigenen Dispatcher herzustellen; geben Sie Err zurück, um den Lauf fehlschlagen zu lassen, wodurch das Wiederholungsbudget sinnvoll wird.

Statuslebenszyklus

queued → running → completed
                 → failed (retries if configured)
                 → cancelled (via DELETE)

Bei einer Wiederholung wird der Workflow von Anfang an mit der ursprünglichen Eingabe erneut ausgeführt. Sie ist nicht checkpoint-bewusst. Daher sollte ein Workflow mit Seiteneffekten idempotent sein oder seinen eigenen Fortschritt absichern.

Wichtig: Laufdatensätze werden in einem In-Memory-Speicher gehalten. Ein Neustart des Prozesses löscht alle Datensätze, einschließlich laufender Läufe. Betrachten Sie diese Endpunkte als asynchrone Ausführungsoberfläche, nicht als dauerhafte Auftragswarteschlange.

Cron-Zeitplanung

Die Funktion background umfasst auch die Verwaltung von Cron-Aufträgen.

Endpunkte

MethodePfadBeschreibung
POST/cronEinen Cron-Job erstellen
GET/cronAlle Jobs auflisten
GET/cron/{job_id}Einen Auftrag abrufen
PATCH/cron/{job_id}Pausieren/fortsetzen
DELETE/cron/{job_id}Einen Auftrag löschen

Parallelitätsrichtlinien

  • skip: Das Vorkommen überspringen, wenn noch eine vorherige Ausführung aktiv ist (Standard)
  • allow: Gleichzeitige Ausführungen zulassen
  • queue: Das Vorkommen in die Warteschlange einreihen und ausführen, sobald die aktive Ausführung abgeschlossen ist

Jedes Vorkommen eines Zeitplans wird einmal beansprucht. Daher erzeugt ein Vorkommen, das eintrifft, während eine Ausführung aktiv ist, genau eine Ausführung in der Warteschlange und nicht eine pro Abfrage des Schedulers. Ausführungen in der Warteschlange werden der Reihe nach abgearbeitet, und jede wird genauso überwacht wie eine direkt ausgelöste Ausführung. Dadurch kehrt die Anzahl aktiver Ausführungen auf null zurück, selbst wenn eine Ausführung fehlschlägt oder abgebrochen wird.

Cron-Aufträge verwenden denselben Executor wie Hintergrundausführungen — ein Auftrag, dessen workflowId nicht registriert ist, führt zum Fehlschlagen seiner Ausführungen.

Verwendung

use adk_server::background::{BackgroundState, CronState, cron_jobs_router_with_state, start_cron_scheduler};

let bg_state = BackgroundState::new();
let cron_state = CronState::new(bg_state);
let app = axum::Router::new().merge(cron_jobs_router_with_state(cron_state.clone()));

// Start the background scheduler
start_cron_scheduler(cron_state);

Beispiele

# Functional API (TaskContext, ReducedValue, MessagesValue, etc.)
cargo run --manifest-path examples/functional_workflow/Cargo.toml

# Background Runs (REST API with Axum)
cargo run --manifest-path examples/background_runs/Cargo.toml

# Cron Scheduling (full lifecycle demo)
cargo run --manifest-path examples/cron_scheduling/Cargo.toml

Feature-Flags

FunktionCrateFügt hinzu
functionaladk-graphTaskContext, typisierte Reducer, Schemavalidierung, Proc-Makros
backgroundadk-serverEndpunkte für Hintergrundausführungen, Cron-Zeitplanung, Scheduler-Schleife

Vorherige: ← Graph Agents | Nächste: Echtzeit-Agenten →

Unterbrechung und typisierter Fortsetzungswert

TaskContext::interrupt<T> unterbricht den Workflow und gibt bei einer späteren Ausführung den für diese Stelle bereitgestellten Wert zurück.

Jede Unterbrechungsstelle erhält anhand ihrer Position in der Ausführung einen Fortsetzungsschlüsselinterrupt-1, interrupt-2 usw. —, sodass ein erneut abgespielter Workflow dieselben Unterbrechungen in derselben Reihenfolge erreicht und den ihm übergebenen Wert vorfindet.

// First run: no value supplied, so the workflow suspends.
let error = ctx.interrupt::<Approval>("approve the refund").await.unwrap_err();
// workflow suspended at interrupt 'interrupt-1': approve the refund

// Second run: supply the value under that key.
let ctx = ctx.with_resume_values(HashMap::from([
    ("interrupt-1".to_string(), serde_json::json!({ "approved": true, "approver": "alice" })),
]));
let approval: Approval = ctx.interrupt("approve the refund").await?;

Der Schlüssel wird außerdem unter continuation_key im Unterbrechungsprüfpunkt gespeichert.

SituationErgebnis
Kein Wert für den SchlüsselFunctionalError::Suspended { continuation_key, message }
Ein Wert, der in T deserialisiert wirdOk(value)
Ein Wert, der nicht in T deserialisiert werden kannFunctionalError::InterruptTypeMismatch zur Benennung der Website

Wichtig: interrupt gab zuvor bei jedem Aufruf InterruptTypeMismatch mit „Workflow unterbrochen“ zurück, und außerhalb der Methode wurde kein Fortsetzungswert verarbeitet. Ein Aufrufer konnte nicht unterscheiden, ob „Eingabe erforderlich“ vorlag oder „Ihr Wert hatte den falschen Typ“, hatte keinen Schlüssel, unter dem ein Wert bereitgestellt werden konnte, und erhielt an der Aufrufstelle nie einen typisierten Wert.

Hinweis: From<FunctionalError> for GraphError wird zu GraphError::Other(String) abgeflacht. Ein Aufrufer, der strukturell abgleicht, sollte daher vor der Konvertierung auf FunctionalError abgleichen.