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
| Methode | Pfad | Beschreibung |
|---|---|---|
| POST | /runs | Einen 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
| Methode | Pfad | Beschreibung |
|---|---|---|
| POST | /cron | Einen Cron-Job erstellen |
| GET | /cron | Alle 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
| Funktion | Crate | Fügt hinzu |
|---|---|---|
functional | adk-graph | TaskContext, typisierte Reducer, Schemavalidierung, Proc-Makros |
background | adk-server | Endpunkte 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üssel — interrupt-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.
| Situation | Ergebnis |
|---|---|
| Kein Wert für den Schlüssel | FunctionalError::Suspended { continuation_key, message } |
Ein Wert, der in T deserialisiert wird | Ok(value) |
Ein Wert, der nicht in T deserialisiert werden kann | FunctionalError::InterruptTypeMismatch zur Benennung der Website |
Wichtig:
interruptgab zuvor bei jedem AufrufInterruptTypeMismatchmit „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 GraphErrorwird zuGraphError::Other(String)abgeflacht. Ein Aufrufer, der strukturell abgleicht, sollte daher vor der Konvertierung aufFunctionalErrorabgleichen.