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étodoCaminhoDescrição
POST/runsEnviar 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étodoCaminhoDescrição
POST/cronCriar um cron job
GET/cronListar 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

RecursoPacoteAdiciona
functionaladk-graphTaskContext, redutores tipados, validação de esquema, macros de procedimento
backgroundadk-serverendpoints 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çãoResultado
Nenhum valor para a chaveFunctionalError::Suspended { continuation_key, message }
Um valor que é desserializado em TOk(value)
Um valor que não é desserializado em TFunctionalError::InterruptTypeMismatch que nomeia o site

Importante: interrupt anteriormente retornava InterruptTypeMismatch com "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 para GraphError::Other(String), portanto um chamador que faça correspondência estrutural deve fazer a correspondência com FunctionalError antes da conversão.