操作节点

adk-action crate 定义了 14 种用于基于图的工作流的操作节点类型。每种节点类型都代表一个离散操作(HTTP 调用、数据转换、条件分支等),可通过 adk-graphActionNodeExecutor 组合成有向图。

概述

操作节点是可视化和编程式工作流图的构建块。它们提供:

  • 类型化操作 — 涵盖常见工作流模式的 14 种节点类型
  • StandardProperties — 用于错误处理、追踪和数据映射的共享配置
  • 变量插值 — 用于动态值的 {{variable}} 语法
  • 图集成 — 由 adk-graphActionNodeExecutor 执行

安装

[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: javascripttypescriptCode已拒绝 — 没有沙箱运行时;请使用 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_errorErrorHandlingStopContinueOnFailRetryThenFail
retryOption<RetryConfig>重试次数及重试之间的延迟
notesOption<String>用于跟踪的可读描述
on_successOption<String>成功时执行的回调节点
on_failureOption<String>失败时执行的回调节点
timeout_msOption<u64>最大执行时间
continue_on_failbool失败后是否执行下游节点
input_mappingOption<String>用于转换输入数据的表达式
output_keyOption<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-graphActionNodeExecutor 执行:

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" }))
    }
}

上一页← 重试与反思 | 下一页插件 →