Funcional API
Escribe flujos de trabajo de agentes como funciones async normales de Rust con checkpointing automático, reductores de estado tipados y compatibilidad con interrupción/reanudación.
Descripción general
El API (funcionalidad functional en adk-graph) proporciona una alternativa de nivel superior a la construcción explícita de grafos/nodos/aristas. En lugar de construir manualmente un StateGraph, anota las funciones con las macros #[entrypoint] y #[task] y usa el flujo de control estándar de Rust.
Beneficios principales:
- Escribe flujos de trabajo como Rust async normal, sin necesidad de un DSL de grafos
- Checkpointing automático después de cada tarea para la recuperación tras fallos
- Flujo de control estándar de Rust (if/else, for, loop, match)
- Contenedores de estado tipados con garantías de persistencia
- Compatible con la infraestructura existente de StreamEvent y Checkpointer
Primeros pasos
[dependencies]
adk-graph = { version = "2.1.0", features = ["functional"] }
adk-rust-macros = "2.1.0"
Tipos principales
TaskContext
El contexto de ejecución que se pasa a todas las funciones del flujo de trabajo. Proporciona acceso al estado, el checkpointing, las interrupciones y el 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>
Contenedor de estado de solo adición que se conserva entre checkpoints. Los valores se acumulan y nunca se sobrescriben.
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>
Valores transitorios de ejecución excluidos de la persistencia de checkpoints. Se restablecen a su valor predeterminado al reanudar.
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
Contenedor de mensajes de chat con deduplicación automática basada en los identificadores de los mensajes.
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
Valida los tipos de estado en los límites del flujo de trabajo y detecta las incompatibilidades con antelación.
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
Registra la finalización de las tareas para permitir omitirlas al reanudar. Las tareas completadas se omiten al reanudar el flujo de trabajo.
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
}
Ejecuciones en segundo plano
La funcionalidad background de adk-server añade endpoints REST para la ejecución asíncrona de flujos de trabajo.
Endpoints
| Método | Ruta | Descripción |
|---|---|---|
| POST | /runs | Enviar una ejecución en segundo plano |
| GET | /runs/{run_id} | Consultar el estado de la ejecución |
| ELIMINAR | /runs/{run_id} | Cancelar una ejecución |
Uso
Una ejecución nombra un workflowId. Algo debe convertir ese nombre en trabajo, por lo que debes registrar un
WorkflowExecutor; sin uno, una ejecución enviada falla en lugar de informar
que se ha completado:
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));
Enviar un workflowId no registrado devuelve 404 en lugar de poner en cola una ejecución que
nunca puede ejecutarse. Implementa WorkflowExecutor directamente para conectarlo con adk-graph, el
API funcional o tu propio distribuidor; devuelve Err para hacer
fallar la ejecución, que es lo que hace significativo el presupuesto de reintentos.
Ciclo de vida del estado
queued → running → completed
→ failed (retries if configured)
→ cancelled (via DELETE)
El reintento vuelve a ejecutar el flujo de trabajo desde el principio con la entrada original. No tiene en cuenta los puntos de control, por lo que un flujo de trabajo con efectos secundarios debe ser idempotente o proteger su propio progreso.
Importante: los registros de las ejecuciones se almacenan en un almacén en memoria. Reiniciar el proceso elimina todos los registros, incluidas las ejecuciones en curso. Trata estos endpoints como una superficie de ejecución asíncrona, no como una cola de trabajos duradera.
Programación con Cron
La funcionalidad de background también incluye la gestión de trabajos cron.
Endpoints
| Método | Ruta | Descripción |
|---|---|---|
| POST | /cron | Crear un trabajo cron |
| GET | /cron | Enumerar todos los trabajos |
| OBTENER | /cron/{job_id} | Obtener un trabajo |
| PARCHEAR | /cron/{job_id} | Pausar/reanudar |
| ELIMINAR | /cron/{job_id} | Eliminar un trabajo |
Políticas de concurrencia
- skip: Omitir la aparición si una ejecución anterior aún está activa (predeterminado)
- allow: Permitir ejecuciones simultáneas
- queue: Poner en cola la aparición y ejecutarla cuando finalice la ejecución activa
Cada aparición programada se reclama una sola vez, por lo que una aparición que llega mientras hay una ejecución activa produce exactamente una ejecución en cola, en lugar de una por cada consulta del programador. Las ejecuciones en cola se procesan en orden y cada una se supervisa de la misma manera que una ejecución activada directamente, por lo que el recuento de ejecuciones activas vuelve a cero incluso si una ejecución falla o se cancela.
Los trabajos de Cron utilizan el mismo ejecutor que las ejecuciones en segundo plano; un trabajo cuyo workflowId no esté registrado provocará que sus ejecuciones fallen.
Uso
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);
Ejemplos
# 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
Indicadores de funcionalidad
| Función | Crate | Añade |
|---|---|---|
functional | adk-graph | TaskContext, reductores tipados, validación de esquemas, macros proc |
background | adk-server | Endpoints de ejecución en segundo plano, programación con cron, bucle del planificador |
Anterior: ← Agentes de grafos | Siguiente: Agentes en tiempo real →
Interrupción y reanudación tipada
TaskContext::interrupt<T> suspende el flujo de trabajo y, en una ejecución posterior, devuelve el valor proporcionado
para ese punto.
Cada punto de interrupción obtiene una clave de continuación a partir de su posición en la ejecución — interrupt-1,
interrupt-2, etc. —, de modo que un flujo de trabajo reproducido alcanza las mismas interrupciones en el mismo orden
y encuentra el valor que se le proporcionó.
// 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 clave también se escribe en el punto de control de la interrupción bajo continuation_key.
| Situación | Resultado |
|---|---|
| No hay ningún valor para la clave | FunctionalError::Suspended { continuation_key, message } |
Un valor que se deserializa en T | Ok(value) |
Un valor que no se deserializa en T | FunctionalError::InterruptTypeMismatch que nombra el sitio |
Importante:
interruptanteriormente devolvíaInterruptTypeMismatchcon «flujo de trabajo interrumpido» en cada llamada, y nada fuera del método consumía un valor de reanudación. La persona que llamaba no podía distinguir entre «se necesita una entrada» y «tu valor era del tipo incorrecto», no tenía ninguna clave con la que proporcionar un valor y nunca recibía un valor tipado en el punto de llamada.
Nota:
From<FunctionalError> for GraphErrorse aplana aGraphError::Other(String), por lo que quien llame y haga una coincidencia estructural debería comparar conFunctionalErrorantes de la conversión.