函数式 API
将 agent 工作流编写为普通的异步 Rust 函数,并获得自动检查点、类型化状态归约器以及中断/恢复支持。
概述
函数式 API(adk-graph 中的 functional 功能)提供了一种比显式构建图/节点/边更高层次的替代方案。无需手动构建 StateGraph,只需使用 #[entrypoint] 和 #[task] 宏标注函数,并使用标准 Rust 控制流。
主要优点:
- 将工作流编写为普通的异步 Rust——无需图 DSL
- 每个任务完成后自动创建检查点,以便从崩溃中恢复
- 标准 Rust 控制流(if/else、for、loop、match)
- 具有持久性保证的类型化状态容器
- 兼容现有的 StreamEvent 和检查点基础设施
快速入门
[dependencies]
adk-graph = { version = "2.1.0", features = ["functional"] }
adk-rust-macros = "2.1.0"
核心类型
TaskContext
传递给所有工作流函数的运行时上下文。提供对状态、检查点、中断和流式传输的访问。
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>
在检查点之间持久化的仅追加状态容器。值会累积,且永不被覆盖。
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>
不纳入检查点持久化的临时运行时值。恢复时会重置为默认值。
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
带有基于消息 ID 自动去重功能的聊天消息容器。
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
在工作流边界验证状态类型——尽早捕获类型不匹配。
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
跟踪任务完成情况,以支持恢复时跳过任务的行为。工作流恢复时会跳过已完成的任务。
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
}
后台运行
adk-server 中的 background 功能为异步工作流执行增加了 REST 端点。
端点
| 方法 | 路径 | 描述 |
|---|---|---|
| POST | /runs | 提交后台运行 |
| GET | /runs/{run_id} | 轮询运行状态 |
| 删除 | /runs/{run_id} | 取消运行 |
用法
一次运行会命名一个 workflowId。必须有某个对象将该名称转换为工作,因此请注册一个
WorkflowExecutor — 如果没有注册,已提交的运行会失败,而不是报告完成:
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));
提交未注册的 workflowId 会返回 404,而不是将其加入队列,导致运行永远无法执行。直接实现
WorkflowExecutor,以连接到 adk-graph、功能型 API 或你自己的调度器;返回 Err
即可使运行失败,这正是重试次数上限发挥作用的原因。
状态生命周期
queued → running → completed
→ failed (retries if configured)
→ cancelled (via DELETE)
重试会使用原始输入从头开始重新执行工作流。它不支持检查点,因此包含副作用的工作流应当具备幂等性,或自行保护其进度。
**重要提示:**运行记录存储在内存存储中。进程重启会丢失所有记录,包括正在执行的运行。请将这些端点视为异步执行接口,而不是持久化任务队列。
Cron 调度
background 功能还包括 Cron 作业管理。
端点
| 方法 | 路径 | 描述 |
|---|---|---|
| POST | /cron | 创建 cron 作业 |
| GET | /cron | 列出所有作业 |
| GET | /cron/{job_id} | 获取一个任务 |
| PATCH | /cron/{job_id} | 暂停/恢复 |
| DELETE | /cron/{job_id} | 删除任务 |
并发策略
- skip:如果之前的运行仍处于活动状态,则跳过此次发生(默认)
- allow:允许并发执行
- queue:将此次发生排入队列,并在活动运行完成后执行
每次计划发生都会被声明一次,因此,在某次运行处于活动状态时到达的发生,只会产生一个排队运行,而不会因调度器轮询而产生多个运行。排队运行会按顺序执行,并且每个运行都会以与直接触发的运行相同的方式受到监控,因此即使某个运行失败或被取消,活动运行计数也会回到零。
Cron 作业使用与后台运行相同的执行器——未注册其 workflowId 的作业,其运行将失败。
用法
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);
示例
# 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
功能标志
| 功能 | Crate | 添加内容 |
|---|---|---|
functional | adk-graph | TaskContext、类型化 reducer、模式验证、过程宏 |
background | adk-server | 后台运行端点、cron 调度、调度器循环 |
中断与类型化恢复
TaskContext::interrupt<T> 会暂停工作流,并在之后的运行中返回为该位置提供的值。
每个中断位置都会根据其在运行中的位置获得一个延续键——interrupt-1、
interrupt-2 等——这样,重放的工作流就会按相同顺序到达相同的中断,并获取提供给它的值。
// 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?;
该键也会写入 continuation_key 下的中断检查点。
| 情况 | 结果 |
|---|---|
| 键没有值 | FunctionalError::Suspended { continuation_key, message } |
反序列化为 T 的值 | Ok(value) |
无法反序列化为 T 的值 | 命名该站点的 FunctionalError::InterruptTypeMismatch |
重要:
interrupt之前在每次调用时都返回InterruptTypeMismatch,并显示“工作流已中断”,且方法外部没有任何部分使用恢复值。调用者无法区分“需要输入”和“你的值类型错误”,没有可用于提供值的键,也无法在调用点获得类型化的值。
注意:
From<FunctionalError> for GraphError会扁平化为GraphError::Other(String),因此进行结构匹配的调用者应在转换前匹配FunctionalError。