操作节点
adk-action crate 定义了 14 种用于基于图的工作流的操作节点类型。每种节点类型都代表一个离散操作(HTTP 调用、数据转换、条件分支等),可通过 adk-graph 的 ActionNodeExecutor 组合成有向图。
概述
操作节点是可视化和编程式工作流图的构建块。它们提供:
- 类型化操作 — 涵盖常见工作流模式的 14 种节点类型
- StandardProperties — 用于错误处理、追踪和数据映射的共享配置
- 变量插值 — 用于动态值的
{{variable}}语法 - 图集成 — 由
adk-graph的ActionNodeExecutor执行
安装
[dependencies]
adk-action = "2.1.0"
# Or specific action features via umbrella crate
adk-rust = { version = "2.1.0", features = ["action"] }
节点类型(14 种)
| 节点 | 用途 | 类别 |
|---|---|---|
Trigger | 入口点 — 启动工作流执行 | 控制 |
HTTP | 发出 HTTP 请求(GET、POST、PUT、DELETE、PATCH) | I/O |
Set | 为工作流变量赋值 | 数据 |
Transform | 使用表达式或代码转换数据 | 数据 |
Switch | 基于表达式进行条件分支 | 控制 |
Loop | 遍历集合或直到满足条件 | 控制 |
Merge | 将多个分支重新合并 | 控制 |
Wait | 暂停执行一段时间或直到事件发生 | 控制 |
Code | 执行任意代码(Rust;JS/TS 尚未实现) | 计算 |
Database | 查询数据库——尚未实现 | 输入/输出 |
Email | 通过 SMTP 发送电子邮件 — 尚未实现 | 输入/输出 |
Notification | 发送通知(Slack、webhook、推送) | 输入/输出 |
RSS | 读取并解析 RSS/Atom 订阅源 | 输入/输出 |
File | 读取、写入和转换文件 | 输入/输出 |
图构建时会检查可用性
某些节点类型会接受并验证配置,但其后端并不存在。
这些节点会被 StateGraph::compile() 拒绝,而不是在运行过程中途拒绝:
| 配置 | 状态 |
|---|---|
Database(任何类型) | 已拒绝 — 未集成驱动程序 |
Email(监控或发送) | 已拒绝 — 未实现 IMAP 和 SMTP |
带有 language: javascript 或 typescript 的 Code | 已拒绝 — 没有沙箱运行时;请使用 rust |
不带 action-http 功能的 Http | 已拒绝 — 启用该功能 |
拒绝会指出节点及其原因,因此无法运行的工作流会在组装过程中失败,而不是在前面的节点已经产生副作用之后才失败。
在自定义节点上实现 Node::validate,即可参与相同的检查。
StandardProperties
每个操作节点都携带 StandardProperties ——用于控制执行行为的共享配置:
use adk_action::{StandardProperties, ErrorHandling, RetryConfig};
let props = StandardProperties::builder()
// Error handling
.on_error(ErrorHandling::ContinueOnFail)
.retry(RetryConfig {
max_attempts: 3,
wait_between_ms: 1000,
})
// Tracing
.notes("Fetch user profile from API")
// Callbacks
.on_success("notify_complete")
.on_failure("alert_team")
// Execution
.timeout_ms(30_000)
.continue_on_fail(true)
// Input/output mapping
.input_mapping("{{trigger.body.user_id}}")
.output_key("user_profile")
.build();
StandardProperties 字段
| 字段 | 类型 | 描述 |
|---|---|---|
on_error | ErrorHandling | Stop、ContinueOnFail 或 RetryThenFail |
retry | Option<RetryConfig> | 重试次数及重试之间的延迟 |
notes | Option<String> | 用于跟踪的可读描述 |
on_success | Option<String> | 成功时执行的回调节点 |
on_failure | Option<String> | 失败时执行的回调节点 |
timeout_ms | Option<u64> | 最大执行时间 |
continue_on_fail | bool | 失败后是否执行下游节点 |
input_mapping | Option<String> | 用于转换输入数据的表达式 |
output_key | Option<String> | 用于存储输出的变量名 |
变量插值
操作节点支持 {{variable}} 语法,用于引用工作流状态:
use adk_action::HttpNode;
let node = HttpNode::builder()
.url("https://api.example.com/users/{{trigger.body.user_id}}")
.method("GET")
.headers(vec![
("Authorization".into(), "Bearer {{env.API_TOKEN}}".into()),
])
.build();
变量来源
| 前缀 | 来源 | 示例 |
|---|---|---|
trigger | 触发器节点负载 | {{trigger.body.email}} |
env | 环境变量 | {{env.DATABASE_URL}} |
nodes | 前序节点的输出 | {{nodes.fetch_user.json.name}} |
workflow | 工作流级变量 | {{workflow.run_id}} |
嵌套访问使用点号表示法:{{nodes.http_1.json.data[0].id}}
节点示例
HTTP 节点
use adk_action::{HttpNode, HttpMethod};
let node = HttpNode::builder()
.method(HttpMethod::Post)
.url("https://api.example.com/orders")
.headers(vec![
("Content-Type".into(), "application/json".into()),
])
.body(r#"{"item": "{{trigger.body.item}}", "qty": {{trigger.body.quantity}}}"#)
.properties(StandardProperties::builder()
.timeout_ms(10_000)
.on_error(ErrorHandling::RetryThenFail)
.retry(RetryConfig { max_attempts: 3, wait_between_ms: 2000 })
.output_key("order_response")
.build())
.build();
Switch 节点
use adk_action::{SwitchNode, SwitchCase};
let node = SwitchNode::builder()
.cases(vec![
SwitchCase {
condition: "{{nodes.classify.json.category}} == 'urgent'".into(),
output: "urgent_path".into(),
},
SwitchCase {
condition: "{{nodes.classify.json.category}} == 'normal'".into(),
output: "normal_path".into(),
},
])
.fallback("default_path")
.build();
循环节点
use adk_action::{LoopNode, LoopMode};
let node = LoopNode::builder()
.mode(LoopMode::ForEach {
items: "{{nodes.fetch_users.json.users}}".into(),
item_var: "current_user".into(),
})
.body_nodes(vec!["process_user", "save_result"])
.properties(StandardProperties::builder()
.notes("Process each user in the list")
.build())
.build();
设置节点
use adk_action::SetNode;
let node = SetNode::builder()
.assignments(vec![
("status".into(), "processing".into()),
("started_at".into(), "{{workflow.timestamp}}".into()),
("user_email".into(), "{{trigger.body.email}}".into()),
])
.build();
与 adk-graph 集成
操作节点由 adk-graph 的 ActionNodeExecutor 执行:
use adk_graph::{Graph, ActionNodeExecutor};
use adk_action::{TriggerNode, HttpNode, SetNode};
// Define nodes
let trigger = TriggerNode::webhook("order_received");
let fetch = HttpNode::get("https://api.example.com/inventory/{{trigger.body.sku}}");
let update = SetNode::new(vec![("available", "{{nodes.fetch.json.quantity}}")]);
// Build graph
let graph = Graph::builder()
.node("trigger", trigger)
.node("check_inventory", fetch)
.node("update_status", update)
.edge("trigger", "check_inventory")
.edge("check_inventory", "update_status")
.build()?;
// Execute
let executor = ActionNodeExecutor::new();
let result = executor.run(graph, initial_context).await?;
定义自定义节点类型
实现 ActionNode trait:
use adk_action::{ActionNode, ActionContext, ActionResult, StandardProperties};
use async_trait::async_trait;
struct CustomNode {
config: MyConfig,
properties: StandardProperties,
}
#[async_trait]
impl ActionNode for CustomNode {
fn node_type(&self) -> &str { "custom" }
fn properties(&self) -> &StandardProperties { &self.properties }
async fn execute(&self, ctx: &ActionContext) -> ActionResult {
let input = ctx.resolve("{{trigger.body.data}}")?;
// Custom logic...
Ok(serde_json::json!({ "result": "processed" }))
}
}
相关内容
- 图代理 — 图工作流编排
- Studio 操作节点 — ADK Studio 中的可视化节点编辑器
- 触发器 — 工作流触发器类型