Funcional API
Escreva fluxos de trabalho de agentes como funções assíncronas Rust normais, com checkpointing automático, redutores de estado tipados e suporte a interrupção/retomada.
Visão geral
O API funcional (recurso functional em adk-graph) oferece uma alternativa de nível mais alto à construção explícita de grafos/nós/arestas. Em vez de criar um StateGraph manualmente, você anota funções com as macros #[entrypoint] e #[task] e usa o fluxo de controle padrão do Rust.
Principais benefícios:
- Escreva fluxos de trabalho como Rust assíncrono normal — sem necessidade de uma DSL de grafos
- Checkpointing automático após cada tarefa para recuperação após falhas
- Fluxo de controle padrão do Rust (if/else, for, loop, match)
- Contêineres de estado tipados com garantias de persistência
- Compatível com a infraestrutura existente de StreamEvent e Checkpointer
Primeiros passos
[dependencies]
adk-graph = { version = "2.1.0", features = ["functional"] }
adk-rust-macros = "2.1.0"
Tipos principais
TaskContext
O contexto de execução passado a todas as funções do fluxo de trabalho. Fornece acesso ao estado, ao checkpointing, às interrupções e ao 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>
Contêiner de estado somente para acréscimo, persistido entre checkpoints. Os valores se acumulam e nunca são substituídos.
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 de execução transitórios excluídos da persistência do checkpoint. São redefinidos para o valor padrão ao retomar.
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
Contêiner de mensagens de chat com deduplicação automática baseada nos IDs das mensagens.
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 os tipos de estado nos limites do fluxo de trabalho — detecta incompatibilidades antecipadamente.
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
Rastreia a conclusão de tarefas para o comportamento de ignorar ao retomar. As tarefas concluídas são ignoradas quando o fluxo de trabalho é retomado.
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
}
Execuções em segundo plano
O recurso background em adk-server adiciona endpoints REST para execução assíncrona de fluxos de trabalho.
Endpoints
| Método | Caminho | Descrição |
|---|---|---|
| POST | /runs | Enviar uma execução em segundo plano |
| GET | /runs/{run_id} | Consultar o status da execução |
| EXCLUIR | /runs/{run_id} | Cancelar uma execução |
Uso
Uma execução nomeia um workflowId. Algo precisa transformar esse nome em trabalho, portanto registre um
WorkflowExecutor — sem um, uma execução enviada falha em vez de informar
a conclusão:
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 um workflowId não registrado retorna 404, em vez de enfileirar uma execução que
nunca poderá ser executada. Implemente WorkflowExecutor diretamente para fazer a ponte com adk-graph, o
API funcional, ou seu próprio despachante; retorne Err para fazer a execução falhar, que é o que torna
o orçamento de tentativas significativo.
Ciclo de vida do status
queued → running → completed
→ failed (retries if configured)
→ cancelled (via DELETE)
A nova tentativa reexecuta o fluxo de trabalho desde o início com a entrada original. Ela não tem consciência de pontos de verificação, portanto um fluxo de trabalho com efeitos colaterais deve ser idempotente ou proteger seu próprio progresso.
Importante: os registros de execução residem em um armazenamento em memória. A reinicialização de um processo perde todos os registros, incluindo as execuções em andamento. Trate esses endpoints como uma superfície de execução assíncrona, não como uma fila de tarefas durável.
Agendamento do Cron
O recurso background também inclui o gerenciamento de tarefas do cron.
Endpoints
| Método | Caminho | Descrição |
|---|---|---|
| POST | /cron | Criar um cron job |
| GET | /cron | Listar todos os jobs |
| GET | /cron/{job_id} | Obter um job |
| PATCH | /cron/{job_id} | Pausar/retomar |
| DELETE | /cron/{job_id} | Excluir um job |
Políticas de concorrência
- skip: Ignore a ocorrência se uma execução anterior ainda estiver ativa (padrão)
- allow: Permita execuções simultâneas
- queue: Coloque a ocorrência na fila e execute-a quando a execução ativa terminar
Cada ocorrência de agendamento é reivindicada uma única vez; portanto, uma ocorrência que chega enquanto uma execução está ativa produz exatamente uma execução na fila, em vez de uma por consulta do agendador. As execuções na fila são processadas em ordem, e cada uma é monitorada da mesma forma que uma execução acionada diretamente, de modo que a contagem de execuções ativas retorne a zero mesmo se uma execução falhar ou for cancelada.
Os trabalhos do Cron usam o mesmo executor que as execuções em segundo plano — um trabalho cujo workflowId não estiver registrado fará com que suas execuções falhem.
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);
Exemplos
# 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
Sinalizadores de recursos
| Recurso | Pacote | Adiciona |
|---|---|---|
functional | adk-graph | TaskContext, redutores tipados, validação de esquema, macros de procedimento |
background | adk-server | endpoints de execução em segundo plano, agendamento cron, loop do agendador |
Anterior: ← Agentes de grafo | Próximo: Agentes em tempo real →
Interrupção e retomada tipada
TaskContext::interrupt<T> suspende o fluxo de trabalho e, em uma execução posterior, retorna o valor fornecido
para esse ponto.
Cada ponto de interrupção recebe uma chave de continuação com base em sua posição na execução — interrupt-1,
interrupt-2 e assim por diante — para que um fluxo de trabalho reproduzido alcance as mesmas interrupções na mesma ordem
e encontre o valor que lhe foi fornecido.
// 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?;
A chave também é gravada no checkpoint da interrupção em continuation_key.
| Situação | Resultado |
|---|---|
| Nenhum valor para a chave | FunctionalError::Suspended { continuation_key, message } |
Um valor que é desserializado em T | Ok(value) |
Um valor que não é desserializado em T | FunctionalError::InterruptTypeMismatch que nomeia o site |
Importante:
interruptanteriormente retornavaInterruptTypeMismatchcom "workflow interrupted" em todas as chamadas, e nada fora do método consumia um valor de retomada. Um chamador não conseguia distinguir entre "precisa de entrada" e "seu valor era do tipo errado", não tinha uma chave para fornecer um valor e nunca recebia um valor tipado no local da chamada.
Observação:
From<FunctionalError> for GraphErroré achatado paraGraphError::Other(String), portanto um chamador que faça correspondência estrutural deve fazer a correspondência comFunctionalErrorantes da conversão.