Fonctionnel API
Écrivez des workflows d’agents sous forme de fonctions Rust asynchrones normales, avec une sauvegarde automatique des points de contrôle, des réducteurs d’état typés et la prise en charge des interruptions et reprises.
Vue d’ensemble
Le API (fonctionnalité functional dans adk-graph) offre une alternative de niveau supérieur à la construction explicite de graphes/nœuds/arêtes. Au lieu de construire manuellement un StateGraph, vous annotez les fonctions avec les macros #[entrypoint] et #[task] et utilisez le flux de contrôle standard de Rust.
Avantages principaux :
- Écrire des workflows en Rust asynchrone normal — aucun DSL de graphe requis
- Sauvegarde automatique d’un point de contrôle après chaque tâche pour la récupération après incident
- Flux de contrôle Rust standard (if/else, for, loop, match)
- Conteneurs d’état typés avec garanties de persistance
- Compatible avec l’infrastructure StreamEvent et Checkpointer existante
Pour commencer
[dependencies]
adk-graph = { version = "2.1.0", features = ["functional"] }
adk-rust-macros = "2.1.0"
Types principaux
TaskContext
Le contexte d’exécution transmis à toutes les fonctions du workflow. Fournit un accès à l’état, aux points de contrôle, aux interruptions et au 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>
Conteneur d’état en ajout seulement, conservé entre les points de contrôle. Les valeurs s’accumulent et ne sont jamais écrasées.
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>
Valeurs d’exécution temporaires exclues de la persistance des points de contrôle. Réinitialisées à leur valeur par défaut lors de la reprise.
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
Conteneur de messages de discussion avec déduplication automatique fondée sur les identifiants des messages.
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
Valide les types d’état aux limites du workflow — détecte rapidement les incompatibilités.
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
Suit l’achèvement des tâches pour le comportement d’ignorance lors de la reprise. Les tâches terminées sont ignorées lors de la reprise du workflow.
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
}
Exécutions en arrière-plan
La fonctionnalité background dans adk-server ajoute des points de terminaison REST pour l’exécution asynchrone des workflows.
Points de terminaison
| Méthode | Chemin | Description |
|---|---|---|
| POST | /runs | Soumettre une exécution en arrière-plan |
| GET | /runs/{run_id} | Interroger l’état de l’exécution |
| SUPPRIMER | /runs/{run_id} | Annuler une exécution |
Utilisation
Une exécution nomme un workflowId. Quelque chose doit transformer ce nom en tâche ; il faut donc enregistrer un
WorkflowExecutor — sans cela, une exécution soumise échoue au lieu de signaler
son achèvement :
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));
La soumission d’un workflowId non enregistré renvoie 404 au lieu de mettre en file une exécution qui
ne pourra jamais être exécutée. Implémentez directement WorkflowExecutor pour faire le lien avec adk-graph, le
API fonctionnel, ou votre propre répartiteur ; renvoyez Err pour faire échouer l’exécution, ce qui rend
le budget de nouvelles tentatives pertinent.
Cycle de vie de l’état
queued → running → completed
→ failed (retries if configured)
→ cancelled (via DELETE)
Une nouvelle tentative réexécute le workflow depuis le début avec l’entrée d’origine. Elle ne tient pas compte des points de contrôle ; un workflow ayant des effets de bord doit donc être idempotent ou protéger lui-même sa progression.
Important : les enregistrements d’exécution résident dans un magasin en mémoire. Le redémarrage d’un processus fait perdre tous les enregistrements, y compris les exécutions en cours. Considérez ces points de terminaison comme une interface d’exécution asynchrone, et non comme une file d’attente de tâches durable.
Planification Cron
La fonctionnalité background inclut également la gestion des tâches cron.
Points de terminaison
| Méthode | Chemin | Description |
|---|---|---|
| POST | /cron | Créer une tâche cron |
| GET | /cron | Lister toutes les tâches |
| GET | /cron/{job_id} | Obtenir une tâche |
| PATCH | /cron/{job_id} | Mettre en pause/reprendre |
| DELETE | /cron/{job_id} | Supprimer une tâche |
Politiques de concurrence
- skip : Ignorer l’occurrence si une exécution précédente est toujours active (par défaut)
- allow : Autoriser les exécutions concurrentes
- queue : Mettre l’occurrence en file d’attente et l’exécuter lorsque l’exécution active est terminée
Chaque occurrence de planification est revendiquée une seule fois ; ainsi, une occurrence qui arrive alors qu’une exécution est active produit exactement une exécution en file d’attente, plutôt qu’une exécution par interrogation du planificateur. Les exécutions en file d’attente sont traitées dans l’ordre, et chacune est surveillée de la même manière qu’une exécution déclenchée directement. Ainsi, le nombre d’exécutions actives revient à zéro même si une exécution échoue ou est annulée.
Les tâches Cron utilisent le même exécuteur que les exécutions en arrière-plan — une tâche dont workflowId n’est pas enregistré échouera lors de ses exécutions.
Utilisation
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);
Exemples
# 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
Indicateurs de fonctionnalité
| Fonctionnalité | Bibliothèque | Ajouts |
|---|---|---|
functional | adk-graph | TaskContext, réducteurs typés, validation de schéma, macros procédurales |
background | adk-server | Points de terminaison d’exécution en arrière-plan, planification cron, boucle du planificateur |
Précédent : ← Agents graphiques | Suivant : Agents en temps réel →
Interruption et reprise typée
TaskContext::interrupt<T> suspend le workflow et, lors d’une exécution ultérieure, renvoie la valeur fournie
pour ce point.
Chaque point d’interruption reçoit une clé de continuation à partir de sa position dans l’exécution — interrupt-1,
interrupt-2, et ainsi de suite — afin qu’un workflow rejoué atteigne les mêmes interruptions dans le même ordre
et retrouve la valeur qui lui a été fournie.
// 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?;
La clé est également écrite dans le point de contrôle de l’interruption sous continuation_key.
| Situation | Résultat |
|---|---|
| Aucune valeur pour la clé | FunctionalError::Suspended { continuation_key, message } |
Une valeur qui se désérialise en T | Ok(value) |
Une valeur qui ne se désérialise pas en T | FunctionalError::InterruptTypeMismatch désignant le site |
Important :
interruptrenvoyait auparavantInterruptTypeMismatchavec « workflow interrupted » à chaque appel, et rien en dehors de la méthode ne consommait de valeur de reprise. L’appelant ne pouvait pas distinguer « entrée requise » de « votre valeur était du mauvais type », ne disposait d’aucune clé sous laquelle fournir une valeur et ne recevait jamais de valeur typée au point d’appel.
Remarque :
From<FunctionalError> for GraphErrors’aplatit enGraphError::Other(String); par conséquent, un appelant effectuant une correspondance structurelle devrait effectuer la correspondance surFunctionalErroravant la conversion.