安装
acts 是一个快速、轻量、可扩展的工作流引擎库,使用 YAML 格式定义工作流并通过消息驱动架构执行和分发消息。
安装 acts 库
通过 cargo 命令安装:
cargo add acts
安装外部存储
外部存储后端(sqlite/postgres/redis/nats/sled)在独立的 acts-store crate 中,
启用对应 feature 后从 acts_store 导入后端,再通过 EngineBuilder::set_store
指定(不设置时默认使用内存存储 MemoryStore):
# SQLite
cargo add acts-store --features sqlite
# PostgreSQL
cargo add acts-store --features postgres
# NATS
cargo add acts-store --features nats
# Redis
cargo add acts-store --features redis
# Sled
cargo add acts-store --features sled
use acts::Engine;
use acts_store::SqliteStore; // 或 PostgresStore / RedisStore / NatsStore / SledStore
use std::sync::Arc;
#[tokio::main]
async fn main() -> acts::Result<()> {
let store = SqliteStore::open("data/acts.db").await?;
let engine = Engine::builder()
.set_store(Arc::new(store))
.start()
.await?;
Ok(())
}
acts 本身只内置 MemoryStore 与自定义存储所需的 KvStore trait;实现
acts::KvStore 的自定义后端同样通过 set_store 注入。
每个数据库只能有一个写入者:引擎的文档锁是进程内的——同一进程内的所有引擎共享
一张锁表(不同数据库用到了相同键时会被一并串行化,只是多了一点争用),但它不
跨进程,所以两个进程写同一个数据库时彼此没有互斥,同一行的并发更新可能让索引行与
数据行不一致:查询命中已不再持有的值,或漏掉当前值。单实例部署不受影响;多实例
部署请各自使用独立数据库,或在引擎之外自行协调(后端条件写,或覆盖读取过程的锁)
——batch 只保证单次写入原子,不代表「读取 + 另一个进程的写入」是原子的。
创建引擎
#![allow(unused)]
fn main() {
use acts::{Engine, Principal};
let engine = Engine::builder().start().await.unwrap();
// executor 代表一个调用者;直接驱动引擎的嵌入式调用者用的是引擎自身的身份
// (见"访问控制"一章)
let executor = engine.executor(&Principal::unrestricted());
}
部署和启动工作流
#![allow(unused)]
fn main() {
use acts::{Engine, Principal, Vars, Workflow};
let engine = Engine::builder().start().await.unwrap();
// 加载 YAML 模型
let model = r#"
id: my_model
name: my model
steps:
- name: step 1
uses: acts.transform.set
params:
a: 10
- name: step 2
uses: acts.transform.code
params: |
return { data: a + 10 };
"#;
let workflow = Workflow::from_yml(model).unwrap();
// 部署模型
let executor = engine.executor(&Principal::unrestricted());
executor.model().deploy(&workflow).expect("fail to deploy workflow");
// 启动工作流
let mut vars = Vars::new();
vars.set("a", 0);
vars.set("pid", "w1");
executor.proc().start(&workflow.id, vars).expect("fail to start workflow");
}
关联项目
| 项目 | 说明 |
|---|---|
| acts-server | 基于 gRPC 的工作流服务 |
| acts-channel | Rust 客户端库 |
| acts-channel-py | Python 客户端库 |
| acts-channel-go | Go 客户端库 |
访问控制
引擎的每一个操作,都要拿执行它的那个调用者身份做一次校验。引擎配置里的 [acl] 段
就是给这些调用者命名的。没有该段时,引擎对任何人开放,但只交出目录:所有请求都
归属到内置的 anonymous 主体,它能 list/get 模型与包,仅此而已——没有别的读取、
没有写入、没有控制类动作、没有管理类动作、没有 snapshot scope、也不能订阅。未配置的
部署是用来“看“已经部署了什么的,不是用来改它、也不是用来读某个流程、某条消息或某个
触发器里的内容。
两种离开该默认的姿态:
- 加上
[acl]——最小可用的一段就是一个token,它让该 token 拥有一切 (等价于 requirepass),并让所有调用者从匿名变为已认证; - 在
[acl]里写enabled = false(显式退出):不做任何校验,所有调用者不受限。 这既是开启 ACL 之前的行为,如今也是刻意的选择,而不再是“没写配置“的结果。 嵌入式调用者用Engine::builder().disable_acl()表达同一件事——测试和本地演示 用的就是这个设置。
任何地方都没有隐式的不受限策略:没写 [acl] 不是,没人认领的流程不是,
没人点名的 snapshot scope 也不是。它们各自对应下面写明的“什么都不读“的情形。
Token 与角色
请求携带 token,token 选中一个角色,角色的 allow / deny
动作模式决定是否放行。deny 优先。没有 token(或 token 未命中任何角色)的请求
默认被拒绝,除非 default_role 指定了兜底角色——想保留只读默认、同时配置其它规则时,
就用 [[acl.role]] name = "anonymous" 加上 default_role = "anonymous"。
Token 按 SHA-256 摘要比较:写 sha256:<64 位十六进制> 可让明文不落配置;
也可以直接写 token 本身,由服务端在加载时哈希。
[acl]
# 无 token / token 未命中角色时的兜底角色;不写则拒绝这类请求
default_role = "guest"
# 简写:一把不受限的 token(等价于 requirepass)
token = "sha256:9f86d081884c7d659a2feaa0c55ad015a3bf4f1b2b0b822cd15d6c15b0f00a08"
[[acl.role]]
name = "operator"
tokens = ["sha256:2c26b46b68ffc68ff99b453c1d30413413422d706483bfa0f98a5e886266e7ae"]
allow = ["model:ls", "model:get", "proc:ls", "proc:get", "task:*", "msg:ls",
"msg:ack", "snap:get", "snap:ls", "acl:whoami"]
deny = ["model:rm", "pack:publish"]
snapshot = { secrets = ["$subject"], profile = ["$subject/*"] }
workdir = "/srv/acts"
[[acl.role]]
name = "guest"
tokens = ["sha256:..."]
allow = ["model:ls", "acl:whoami"]
配置错误一律是启动错误,绝不静默放行或静默拒绝:模式编译失败、同一 token
被分配给两个角色、default_role 指向不存在的角色、或已启用但没有任何角色声明 token。
动作名
allow/deny 匹配共享动作表里的动作名,支持 *、? 通配——例如 act:*
覆盖全部 act 操作。这些名字与 CLI、channel 客户端使用的完全一致。
| 分类 | 动作 |
|---|---|
| 读取 | model:ls model:get proc:ls proc:get task:ls task:get msg:ls msg:get evt:ls evt:get pack:ls pack:get snap:get snap:ls |
| 写入 | model:deploy pack:publish snap:upsert snap:remove |
| 控制 | proc:start proc:start_from_model act:push act:remove act:submit act:complete act:abort act:cancel act:back act:skip act:error evt:start msg:ack |
| 订阅 | msg:sub |
| 管理 | model:rm pack:rm msg:rm msg:clear msg:redo msg:unsub |
| 嵌入式 | ext:register_var(ext().register_package 走 pack:publish) |
最后一行是嵌入式调用者自己的接口——把用户变量模块装进表达式环境、发布包定义。 线上动作表里没有任何一项通到它,因此传输层的调用者够不着;扩展自己所托管的引擎的 嵌入式调用者,传的是它自己的身份(见executor)。
allow = ["*"] 表示不受限:所有动作放行,所有 snapshot scope 也放行。
acl:whoami 返回调用者自身的身份与生效模式;对已认证调用者隐式放行,
因此可以作为启动自检使用,而不会额外放开任何权限。
没有 [acl] 段时归属的 anonymous 主体,拿到的正是
model:ls model:get pack:ls pack:get——目录。这是引擎无法指名道姓的调用者
可以被授予的最小集合,而且这份清单由测试钉死,不留给解读。其余全部不在其中,
包括命名调用者看来理所当然的读取:一条流程行指名了谁跑了什么,一条投递指名了它
本是发给谁的,一个触发器指名了它将启动哪个模型——引擎认不出的调用者,不是这些行
所描述的那个人。msg:sub 同理,还多一条理由:一条订阅流既承载实时载荷,又会为它
持有的通道逐条写入投递行;snap:get/snap:ls 也不在其中,因为 snapshot scope
只有在策略点名时才有归属。
executor
引擎的操作收在同一个对象上,即 executor,它每个方法都在运行前先做一次校验——
model().deploy()、proc().start() 以及其余全部。创建时就绑定了调用者:
#![allow(unused)]
fn main() {
// 传输层的请求:token 解析完之后
let executor = engine.executor(&principal);
executor.proc().start("my_model", vars).await?;
// 没有携带 token 的请求
let executor = engine.executor(&engine.anonymous());
// 引擎自身的操作;测试与本地演示也传这个
let executor = engine.executor(&Principal::unrestricted());
}
executor 还决定它启动的流程带上什么:proc().start() 与 evt().start()
把该主体的 snapshot scope 与 workdir 根封进流程,因此调用者无法通过在请求里塞一个
权威来放宽自己的读取范围。同一个模型的两个调用者,各自读到的是自己拥有的数据。
嵌入式调用者不在其外:它同样通过 executor 抵达引擎,因此它的操作按它传入的主体
校验。传 Principal::unrestricted() 是一句声明(“这是引擎自身的操作”,或“这个部署
主动退出“),而且写在了调用处。
Snapshot 的 scope 归属
Snapshot target 由 target × scope 寻址(见快照密封数据)。
角色的 snapshot 表限定该角色拥有哪些 target 的哪些 scope;$subject
会替换为角色名,所以 ["$subject"] 表示“只读我自己的 scope“。
该规则在两处强制:
snap:*动作上——读写超出调用者集合的 scope 会被拒绝,snap:ls只返回该主体拥有的 scope,一个租户无法枚举他人的数据。- 密封时——流程带上启动它的那份规则,调度器在把 snapshot 值冻结进任务前会再校验
一次。因此即使有人用别人的
uid启动流程,工作流也读不到另一个主体已密封的数据。
流程的权威只有一个来源:用哪个主体创建 executor 启动它。子流程继承父流程的权威,
因此永远不可能读到开启它的那个流程读不到的东西。完全没有调用者的启动——schedule
触发器、嵌入式调用者直接调 Runtime::start——不带任何权威,也就读不到任何
snapshot scope:权威缺失不等于权威无限,需要受管数据的任务会在密封时带着它缺的那个
主体报错,而不是被塞给整个数据面。
消息面
消息与投递同其它操作一样,由动作授权决定——发出它的流程归谁所有,并不影响 谁能读它、谁能 ack 它。
- 订阅是一个动作。 打开一条流(gRPC
on_message、SSE/msg/sse)需要角色 的allow里含msg:sub;不含时传输层直接回PERMISSION_DENIED/403, 不发流。传输层把收到的 client id 交给该动作,并用动作返回的键注册通道,因此 “被校验的路径“与“实际占用的键“不可能各走各的。 - 键按主体划分命名空间:
{subject}/{传输自身的 id}——主体作前缀,后面接传输 自己的 client id(SSE 保留自己的acts-flow-client-段)。因此第二个调用者用 别的主体已经用过的 client id 订阅时,得到的是另一个通道,而不是顶替那个主体 的 handler。msg:unsub由同一个 id 组合出同一个键,所以调用者只能取消自己命名 空间里的通道;角色名含/会在加载配置时被拒绝,前缀的无歧义正是靠这一点保证。 - 投递只看过滤器与授权:通道会收到所有匹配它自报的
type/state/uses/optionsglob 的消息,无论发出它的流程是谁启动的;拿到消息后能做什么,由调用者 持有的动作决定。不希望读到消息载荷的角色,就是没有msg:sub(以及没有msg:ls/msg:get)的角色。 - 订阅的积压是有界的。 每个订阅者只有一个定长队列(
[grpc].queue_size默认 128,[web].queue_size默认 100):投递不等待客户端,队列满即表示客户端 已经停止读取,该订阅随即被断开——引擎不会为每条塞不进去的消息留一个等待发送的 任务。引擎仍欠这个通道的消息不会随断开而消失:发出它的流程尚未结束、且尚未 ack 的投递,会在客户端以同一个 id 重新订阅(落在同一个通道键上)后被重试定时器 重新投递;与任何一次断线一样,通道只会收到它注册期间发出的消息。 msg:ack与msg:unsub同样是普通动作:拿到授权就能 ack 任意投递 id、取消自己 命名空间内的任意通道。投递 id 不是按客户端寻址的,因此msg:ack应按“写入“级别 授予——它的持有者可以压掉别的调用者尚未 ack 的投递。
各层各管什么
四类规则,各自在能表达它的那一层校验:
| 规则 | 作用对象 | 依据 |
|---|---|---|
| 动作权限 | 每一个操作,同一张表 | 角色的 allow/deny 模式 |
| snapshot scope 归属 | snap:* 动作,以及密封时再校验一次 | 角色的 snapshot 表($subject) |
| 通道命名空间 | 通道键与 msg:unsub | 已认证主体 |
| 目录限定 | 流程的 workdir | [acl]/角色的 workdir |
每一个操作都走动作表,这正是校验得以普遍的原因:传输层把 token 解析成主体, 嵌入式调用者同样如此——操作所运行的 executor 带着那个主体,而动作表与 executor 匹配的是同一批动作名,每个操作只定义一处。唯二没有身份可言的调用者,动作表 也为它们各自备好了答案:
- 不带 token 的请求落到
default_role;引擎没有[acl]段时落到只读的anonymous主体——它仍然是一个调用者,因此照样被校验; - 没人认领的启动(
schedule触发器)与子流程(继承父流程的权威)是引擎自身的启动, 没有调用者可校验:它们能读什么,取决于最终带上的那份权威,而两者都无法从外面 获得权威。
目录控制
配置里写的是根目录:workdir 可写在 [acl] 上,也可写在单个角色上覆盖。策略启动的
每个流程得到的是它下面的自己的目录 <workdir>/<pid>,文件系统访问被限定在那里;进程 id
成为路径的一段——因此不能作为单个安全目录名的 pid(空、.、..,或含路径分隔符、冒号)
会被拒绝,而不是被放到根目录之外。
根目录走与 scope 权威同一条私有通道——它本身就是权威的一部分:随 owner 权威封入进程
env($env 代理拒绝该私有键),随进程持久化,且不作为启动参数下发,因此调用者无法指定
流程被限定到哪个目录。流程自己的目录(即 <root>/<pid>)由三处读到:act 用
Context::workdir(),工作流脚本用 $env.WORK_DIR,两者指的是同一个目录,且都不是配置
里写的那个根。$env.WORK_DIR 由引擎应答,写入会被丢弃,脚本无法改写自己所在的目录。
目录的生命周期与进程的持久行一致:进程结束且消息投递结清后,sweeper 删除行时一并 删除目录(投递出错、等待人工重投的流程保留行,目录也随之保留);从未落盘的启动则 立即删掉自己的目录。因此留在目录里的文件不会比流程本身活得更久。
acts.app.shell 使用它:脚本以该目录为工作目录运行,HOME、TMPDIR/TEMP/TMP、
PWD 都指向其中,ACTS_WORKDIR 让脚本能直接引用自己的目录。脚本中出现绝对路径
(/etc/passwd、C:\Windows)或 .. 段时,会在运行前被拒绝。
该包另有一层 [shell] 配置:allow/deny 两张 glob 清单,匹配整段脚本文本
(* 匹配任意字符,含 / 与换行),deny 优先;两张都为空表示不限制。不合法的
模式是启动错误,绝不静默放行。与目录检查一样,它是策略而非沙箱:文本 glob 看不到
脚本将做什么(a=rm; $a -rf / 里没有一个被禁的词),因此它用于把意图写明、把明显
的那一类拒掉,真正的边界仍然要靠操作系统。
不配置 workdir 时不做任何目录控制,进程可以访问服务端账号能访问的一切,
即该选项出现之前的行为。
各传输的凭证
| 传输 | token 的携带方式 |
|---|---|
| gRPC | 每次请求的 authorization: Bearer <token> metadata,包括 on_message 订阅 |
| HTTP | authorization: Bearer <token> 请求头;/health 保持开放以便探活 |
| NATS | 动作 JSON body 里的 token 字段——broker 认证的是连接,不是单次请求 |
客户端
# CLI:命令行参数优先于环境变量
acts-cli --token "$TOKEN"
ACTS_TOKEN="$TOKEN" acts-cli
CLI 在进入 REPL 前先用 acl:whoami 确认身份,因此 token 缺失或过期会在启动时
暴露,而不是拖到第一条命令。
#![allow(unused)]
fn main() {
use acts_channel::ActsChannel;
let mut client = ActsChannel::connect_with_token("http://127.0.0.1:10080", Some(token)).await?;
}
connect 是不带 token 的形式,即匿名调用者:对没有 [acl] 段的服务端只能读,
对已配置的服务端则会被拒绝,除非 default_role 接纳它。
客户端 Channel
客户端 Channel 用于在应用服务连接工作流服务,用以订阅消息、执行工作流任务等工作。
安装
通过 cargo 命令安装客户端库:
cargo add acts-channel
支持的语言
| 语言 | 库 |
|---|---|
| Rust | acts-channel |
| Python | acts-channel-py |
| Go | acts-channel-go |
基本用法
#![allow(unused)]
fn main() {
use acts_channel::{Client, ChannelOptions};
let mut client = Client::new("http://localhost:8080", &ChannelOptions::default());
// 连接服务
client.connect().await?;
// 订阅消息
client.subscribe("client1", "act*", None, None).await?;
// 部署模型
client.deploy(&model_str).await?;
// 启动工作流
let mut vars = Vars::new();
vars.set("input", 100);
client.start("model_id", vars).await?;
}
安装 acts-channel client库
通过 cargo 命令安装:
cargo add acts-channel
连接
应用服务通过 ActsChannel 创建并连接到服务。
#![allow(unused)]
fn main() {
use acts_channel::ActsChannel;
let mut client = ActsChannel::new("http://localhost:8080");
// 连接到 acts-server
client.connect().await?;
}
订阅
通过客户端 Channel 订阅工作流消息。
订阅消息
#![allow(unused)]
fn main() {
use acts_channel::{ActsChannel, ActsOptions};
let mut client = ActsChannel::connect("http://127.0.0.1:10080").await?;
// ActsOptions 的属性支持 glob 模式,如 "act*" 匹配所有以 act 开头的消息
let options = ActsOptions {
state: Some("{created,completed}".to_string()),
r#type: Some("act*".to_string()),
// 其他配置
..ActsOptions::default()
};
let sub = client
.subscribe(
"client-1",
move |message| {
println!("{message:?}");
},
// 不结束订阅的故障:负载解码失败、自动 ack 失败
move |err| eprintln!("subscription fault: {err}"),
&options,
)
.await?;
// 订阅结束:服务端正常关闭为 Ok(()),流或连接失败为 Err(status)
if let Err(err) = sub.wait().await {
eprintln!("subscription closed: {err}");
}
}
该 id 在服务端是调用者主体命名空间下的一段({subject}/{client_id}):不同主体的
两个订阅者都可以叫 client-1 而互不冲突;在 [acl] 下,订阅只承载该主体启动的流程
的消息,而不是所有租户的。见访问控制。
消息类型
| 类型 | 说明 |
|---|---|
workflow | 流程级别消息 |
step | 步骤级别消息 |
act | 活动消息 |
部署
通过客户端 Channel 部署工作流模型。
部署模型
#![allow(unused)]
fn main() {
use acts_channel::ActsChannel;
let mut client = ActsChannel::connect("http://localhost:8080");
// 从文件加载模型并部署
let model = std::fs::read_to_string("workflow.yml").unwrap();
let resp = client
.deploy(yml, Some("custom_model_id")).await?;
}
动态构建模型
也可以通过 Rust 代码动态构建模型后部署:
#![allow(unused)]
fn main() {
use acts::{Workflow, Vars};
let workflow = Workflow::new("my_model", "my workflow")
.with_step(|step| {
step.with_uses("acts.core.irq", Vars::new().with("key", "my_key"))
});
let model_str = serde_yaml::to_string(&workflow).unwrap();
client.deploy(&model_str, Some("custom_model_id")).await?;
}
启动
通过客户端 Channel 启动工作流。
启动工作流
#![allow(unused)]
fn main() {
use acts_channel::{ActsChannel, ChannelOptions};
let mut client = ActsChannel::connect("http://localhost:8080");
// 启动工作流
let mut vars = Vars::new();
vars.set("a", 100);
client.start("model_id", vars).await?;
}
启动参数
启动时可以传递变量来覆盖工作流的默认 vars:
#![allow(unused)]
fn main() {
let mut vars = Vars::new();
vars.set("input_value", 42);
vars.set("user_name", "admin");
client.start("my_workflow", vars).await?;
}
执行
通过客户端 Channel 对活动执行操作。
完成活动
#![allow(unused)]
fn main() {
let mut options = Vars::new();
options.set("result", "done");
client.complete(&pid, &tid, options).unwrap();
}
触发错误
#![allow(unused)]
fn main() {
let mut options = Vars::new();
options.set("ecode", "err_custom");
client.fail(&pid, &tid, options).unwrap();
}
回退到指定步骤
#![allow(unused)]
fn main() {
let mut options = Vars::new();
options.set("to", "step1");
client.back(&pid, &tid, options).unwrap();
}
取消活动
#![allow(unused)]
fn main() {
let mut options = Vars::new();
options.set("to", "step1");
client.cancel(&pid, &tid, options).unwrap();
}
跳过活动
#![allow(unused)]
fn main() {
client.skip(&pid, &tid, Vars::new()).unwrap();
}
中止活动
#![allow(unused)]
fn main() {
let mut options = Vars::new();
options.set("uid", "u1");
client.abort(&pid, &tid, EventAction::Abort, options).unwrap();
}
移除活动
#![allow(unused)]
fn main() {
client.remove(&pid, &tid, EventAction::Remove, Vars::new()).unwrap();
}
命令行工具
用户可通过内置命令行工具查看数据,启动命令如下:
acts-cli -h <服务端ip> -p <端口>
查看帮助
> help
订阅
客户端可通过sub命令订阅消息
sub <client_id> [type] [state] [tag] [key]
subscribe server message
type, state and tag are all support glob string
client_id: client id
type: message types are in workflow, job, step, branch and act.
state: message state in created, completed, error, cancelled, aborted, skipped and backed.
tag: message tag which is defined in workflow model tag attribute.
key: message key
for examples:
1. sub all messages:
sub 1
2. sub all act messages:
sub 1 act
3. sub created and complete messages
sub 1 * {created,completed}
4. sub all messages that the tag starts with abc
sub 1 * * abc*
5. sub all messages that the key starts with 123
sub 1 * * * 123*
部署
部署流程模型文件
deploy <path>
deploy a workflow
path: yml model local file path
启动
启动流程
start <mid>
start a workflow
mid: workflow model id
管理
模型列表
列出已部署的模型
models [count]
query the current deployed models
count: expect to load the max model count
查看模型
查看模型数据
model <mid> [fmt]
query the model data
mid: model id
fmt: display format with text|json|tree
流程列表
列出运行中的所有流程
procs [count]
query the current running procs
count: expect to load the max proc count
查看流程
查看流程数据
proc <pid> [fmt]
query the proc data
fmt: display format with json|tree, the default is tree
任务列表
列出流程所有任务列表
tasks <pid>
query the proc tasks
pid: the proc id
查看任务
查看任务数据
task <pid> <tid>
query the task data
执行
执行活动针对的是类型为req的活动, 当服务端生成req活动后,该活动处于中断状态,等待客户端执行。
env
每一个命令执行, 都需要一些options参数,该命令生成后续执行动作所需的options参数。
env <op> [key] [value] [value-type]
op: command with set, get, ls
set: set key and value.
get: get by key name
ls: list all env values
json: show in json format
key: env key with string type
value: env value
value-type: value type with string, int, float and json, the default type is string
增加
增加一个请求(req)活动
push <pid> <tid>
push an action to a step
pid: proc id
tid: step task id
extra options:
id: act id, it is reqiured
name: act name
inputs: input parameters
outputs: expose vars to its parents
rets: limits the request options when acting
删除
删除一个活动
remove <pid> <tid>
remove an action
pid: proc id
tid: task id
提交
提交一个活动
submit <pid> <tid>
submit an action
pid: proc id
tid: task id
完成
完成一个活动
complete <pid> <tid>
complete the action
pid: proc id
tid: task id
退回
退回一个活动
back <pid> <tid>
back to the history task
pid: proc id
tid: task id
options:
to: set a step id to point out which step to back
撤销
撤销一个已完成,但是下一步骤还没有完成的活动
cancel <pid> <tid>
cancel the act that is completed before
pid: proc id
tid: task id
跳过
跳过一个活动,断续执行下一步
skip <pid> <tid>
skip the action
pid: proc id
tid: task id
终止
终止一个活动,然后整个流程终止
abort <pid> <tid>
abort the workflow
pid: proc id
tid: task id
错误
将一个活动设为错误,如果该活动没有设置异常处理, 则错误逐级传递,直到整个流程结束。
error <pid> <tid>
set an action as error
pid: proc id
tid: task id
options:
err_code: error code, it is required
err_message: error message
流程模型
工作流引擎的执行依赖于流程模型,acts 流程模型是一个规范化的 YAML 格式文件。
模型结构
一个完整的流程模型由以下部分组成:
id: my_model
name: 模型名称
# 默认变量
vars:
- name: value
value: 0
# 输入 schema (JSON Schema)
inputs:
type: object
properties:
value:
type: number
# 输出 schema (JSON Schema)
outputs:
type: object
properties:
data:
type: object
# 启动触发器
on:
- id: event1
kind: manual
# 执行选项
options:
exposes:
- name: output_key
# 步骤列表
steps:
- id: step1
uses: acts.core.irq
核心概念
| 概念 | 说明 |
|---|---|
| 步骤 | 流程的基本执行单元,使用 uses 指定包 |
| 分支 | 条件分支,通过 if 条件决定执行路径 |
| 活动 | 实际动作执行体,使用 uses 指定功能包 |
| 配置 | 工作流的全局配置,包括变量、事件、输入输出 |
| 包 | 可复用的功能模块,分为 core、transform、event 三类 |
输入
工作流的 inputs 定义了流程模型的输入 JSON Schema,用于在启动工作流时验证输入数据。
id: m1
name: test
inputs:
type: object
properties:
a:
type: integer
default: 5
b:
type: string
default: abc
启动时动态修改
工作流的 inputs 可以在启动时通过 vars 参数进行动态赋值:
#![allow(unused)]
fn main() {
let mut vars = Vars::new();
vars.set("a", 100);
vars.set("b", "new_value");
executor.proc().start(&workflow.id, vars)?;
}
步骤输入
步骤也可以定义 vars 作为本地变量:
steps:
- id: step1
vars:
- name: local_var
value: 10
uses: acts.core.irq
params:
key: act1
输出
导出变量
使用 exposes 控制哪些变量在完成时被导出:
id: m1
name: test
exposes:
- name: a
- name: result
steps:
- id: step1
uses: acts.transform.code
params: |
let a = $get("a");
$set("result", a * 2);
exposes 项省略 type 时不会默认当作 string:导出时按变量在运行时的实际类型
校验并导出。提供字面量 value 时按其值推断类型(如 value: 10 即 number);
显式声明 type 时仍按声明类型严格校验。
步骤导出
步骤也可以通过 options.exposes 导出变量到父级(工作流或上层步骤):
steps:
- id: step1
uses: acts.core.irq
params:
key: act1
exposes:
- name: v
触发器(Triggers)
工作流通过 on 字段声明触发器(Trigger),当触发器被触发时启动工作流。触发器只描述工作流的启动方式,本身不会在流程内执行。
支持的触发器类型(kind):
| kind | 说明 | 触发方式 |
|---|---|---|
manual | 手动触发 | executor.evt().start("model-id:trigger-id", &vars).await,立即返回进程 id |
chat | 聊天触发 | 同上,以字符串消息作为输入(写入 params 变量) |
hook | 钩子触发 | 同上,阻塞等待工作流完成并返回其输出 |
schedule | 定时触发 | 引擎定时器按 cron 表达式自动触发,不可手动启动 |
id: m1
name: test
on:
- id: event_manual
kind: manual
name: start by manual
# 触发时无调用方参数时使用的默认输入
params:
value: 0
- id: event_hook
kind: hook
- id: event_chat
kind: chat
# cron 表达式 6 段: 秒 分 时 日 月 周
- id: event_schedule
kind: schedule
schedule: "0 * * * * *"
params:
value: 0
manual/chat/hook通过executor.evt().start("model-id:trigger-id", &payload).await触发;payload为空时使用声明里的params作为启动输入。manual触发器也可以作为 web url trigger:HTTP 传输层(如acts-plugin-web的POST /hooks/{model-id}:{trigger-id})以请求体为 payload 启动它,因此不需要单独的webhookkind。schedule触发器在部署时记录运行状态(last_run/next_run),由引擎定时器轮询触发;重新部署模型时引擎会同步触发器数据——声明变化则更新、被移除的触发器会被清除。kind也可以填其它已注册事件包 id,走包注册表触发(自定义触发器)。
异常
异常处理机制允许在步骤发生错误时进行恢复处理。
步骤级异常
步骤通过 catches 定义异常处理,catches 是 Step 列表:
steps:
- id: step1
uses: acts.core.irq
params:
key: act1
catches:
# 匹配特定错误码
- uses: acts.core.msg
if: $ecode() == 'err1'
params:
key: catch_err1
# 匹配所有未处理的错误
- uses: acts.core.msg
params:
key: catch_others
错误处理流程
- 步骤中的活动触发错误(通过
EventAction::Error或acts.core.action) - 引擎依次检查
catches列表中的条件 - 匹配到第一个满足条件的 catch 后执行对应的处理步骤
- 处理后步骤继续正常执行
- 如果没有匹配的 catch,错误向上传递
错误码
错误码通过 ecode 传递:
#![allow(unused)]
fn main() {
let mut options = Vars::new();
options.set("ecode", "err1");
rt.do_action2(&pid, &tid, EventAction::Error, options).unwrap();
}
在 catch 的 if 条件中使用 $ecode() 获取错误码。
包
包(Package)是工作流中可复用的执行单元。每个步骤或活动通过 uses 指定包名来调用对应的包。
内置包
| 包名 | 类型 | 说明 |
|---|---|---|
acts.core.irq | 中断 | 发起中断请求,等待客户端处理完成 |
acts.core.msg | 消息 | 发送单向消息到客户端 |
acts.core.block | 块 | 包含子活动列表,支持 sequence 模式 |
acts.core.parallel | 并行 | 对集合进行并行执行 |
acts.core.sequence | 顺序 | 对集合进行顺序执行 |
acts.core.subflow | 子流程 | 调用另一个工作流模型 |
acts.core.action | 命令 | 执行引擎命令 |
acts.transform.set | 变换 | 设置变量值 |
acts.transform.code | 变换 | 执行 JavaScript 代码 (QuickJS 运行时) |
工作流的启动触发器通过 on 字段声明(manual/chat/hook/schedule),参见 触发器。
使用示例
steps:
- id: step1
uses: acts.core.irq
params:
key: act1
- id: step2
uses: acts.core.block
params:
mode: sequence
acts:
- uses: acts.core.irq
params:
key: sub_act1
- uses: acts.core.msg
params:
key: sub_msg
自定义包
可以通过实现 ActPackage trait 来扩展自定义包,用
engine.executor(&principal).ext().register_package(&meta) 注册(principal 是调用者
身份,见访问控制;引擎自身在启动时注册内置包也走同一条路)。
步骤
步骤(Step)是流程的基本执行单元。每个步骤可以使用一个内置或自定义包(uses),并传递参数(params)。步骤按顺序执行,也可以包含分支、异常处理和超时处理。
name: test
steps:
- id: step1
name: step 1
uses: acts.core.irq
params:
key: act1
- id: step2
name: step 2
uses: acts.transform.set
params:
a: 10
步骤包含的属性有:
| key | 名 称 | 说 明 |
|---|---|---|
| id | 标识 | 步骤节点的唯一标识 |
| name | 名称 | 有意义的名称,可以是中文等任意字符 |
| desc | 描述 | 步骤描述信息 |
| tag | 标签 | 标签设置 |
| rn | 资源名 | 用于权限控制的资源名称 |
| uses | 包名 | 使用的包名称,如 acts.core.irq、acts.transform.set 等 |
| params | 参数 | 传递给包的参数 |
| vars | 变量 | 本地变量定义 |
| if | 条件 | 根据条件判断是否跳过当前步骤执行,如 ${{ a }} > 0 |
| while | 循环条件 | 循环:条件满足时反复执行当前步骤 |
| catches | 异常 | 步骤错误后进行异常处理,类型为 Step 列表 |
| timeouts | 超时 | 超时处理,类型为 Step 列表 |
| branches | 分支 | 步骤分支,一个步骤可以有多个分支 |
| next | 下一步 | 当步骤完成后,可直接跳转到指定步骤执行 |
| options | 选项 | 额外选项,如 exposes 导出变量 |
| metadata | 元数据 | 用于UI样式的额外信息,不发送给客户端 |
带 while 的步骤是有界循环:每轮迭代前都会重新求值条件,满足则重复执行;
条件不满足时该步骤被跳过,流程按声明顺序落到它后面的步骤:
steps:
- id: add
while: index < input
uses: acts.transform.code
params: |
$set("value", value + index);
$set("index", index + 1);
- id: end
if 条件不满足时步骤同样被跳过并落到后面的步骤(此时不会执行指向自身或向后的
next),因此 if 与 next 保持原有语义,while 与 next 不能同时使用。
多活动步骤
当步骤需要执行多个活动时,可以使用 acts.core.block 包,在 params 中嵌套定义子活动列表:
steps:
- id: step1
uses: acts.core.block
params:
mode: sequence
acts:
- uses: acts.core.irq
params:
key: act1
- uses: acts.core.msg
params:
key: msg1
acts.core.parallel 可以对集合进行并行执行,acts.core.sequence 可以对集合进行顺序执行。
配置
步骤的 vars 用于定义步骤级别的本地变量:
name: test
steps:
- id: step1
vars:
- name: local_a
value: 5
uses: acts.core.irq
params:
key: act1
步骤也可以通过 options.exposes 导出变量到父节点:
name: test
steps:
- id: step1
uses: acts.core.irq
params:
key: act1
options:
exposes:
- name: a
异常
当步骤发生异常错误时,可以通过 catches 定义异常处理。catches 是一个 Step 列表,每个 catch 步骤通过 if 条件匹配对应的错误。
name: test
steps:
- id: step1
uses: acts.core.irq
params:
key: act1
catches:
# 匹配错误码 err1
- uses: acts.core.msg
if: $ecode() == 'err1'
params:
key: catch1
# 匹配错误码 err2
- uses: acts.core.msg
if: $ecode() == 'err2'
params:
key: catch2
# 捕获所有其他错误(无条件)
- uses: acts.core.msg
params:
key: others
Catch 步骤的属性与普通 Step 相同,可以使用 uses 和 params。if 条件中使用 $ecode() 获取错误码。当 catch 处理后,步骤继续正常执行;如果没有匹配的 catch,错误会向上传递。
超时
当步骤需要检测超时时,可以通过 timeouts 定义超时处理。timeouts 是一个 Step 列表,每个 timeout 步骤通过 if 条件中的 $cost_in() 函数匹配时间阈值。
$cost_in() 支持的时间单位:
s— 秒 (seconds)m— 分钟 (minutes)h— 小时 (hours)d— 天 (days)
name: test
steps:
- id: step1
uses: acts.core.irq
params:
key: act1
timeouts:
# 超时在>=2s 且 <5s - 发送消息
- uses: acts.core.msg
if: $cost_in('2s', '5s')
params:
key: step1_timeout_2s
# 超时在>=5s 且 <8s - 发起中断请求
- uses: acts.core.irq
if: $cost_in('5s', '8s')
params:
key: step1_timeout_5s
# 8秒后超时 - 触发错误
- uses: acts.core.action
if: $cost_in('8s')
params:
action: error
options:
ecode: err_timeout_8s
使用超时需要配置 tick 间隔:
#![allow(unused)]
fn main() {
let engine = Engine::builder()
.tick_interval_secs(1)
.start()
.unwrap();
}
步骤触发器
步骤没有独立的触发器声明。工作流的启动通过 on 字段声明触发器,参见 触发器。
步骤的状态变化可以通过客户端 Channel 的 on_message 回调进行监听。当步骤启动时,引擎会发送步骤创建消息(e.is_type("step")),当步骤完成时,会发送步骤完成消息。
条件
步骤中的条件 if 用来判断执行时是否可跳过当前步骤。
name: test
steps:
- id: step1
if: '${{ a }} > 0'
uses: acts.core.irq
params:
key: act1
- id: step2
uses: acts.core.msg
params:
key: done
条件表达式中使用 ${{ var }} 语法引用变量。当 if 条件为 false 时,步骤会被跳过。
活动列表
一个步骤可以直接指定 uses 和 params 来执行单个活动。当需要执行多个活动时,使用 acts.core.block 包在 params 中嵌套定义子活动列表。
单个活动
steps:
- id: step1
uses: acts.core.irq
params:
key: act1
多个活动(顺序执行)
steps:
- id: step1
uses: acts.core.block
params:
mode: sequence
acts:
- uses: acts.transform.set
params:
a: 10
list:
- u1
- u2
- uses: acts.core.irq
params:
key: act1
- uses: acts.core.msg
params:
key: msg1
并行执行
steps:
- id: step1
uses: acts.core.parallel
params:
in: '${{ list }}'
acts:
- uses: acts.core.irq
params:
key: act2
顺序执行
steps:
- id: step1
uses: acts.core.sequence
params:
in: '${{ list }}'
acts:
- uses: acts.core.irq
params:
key: act2
更多活动内容请参见 活动。
分支
一个步骤可以有多个分支。第一个分支需要设置分支条件,当分支执行时,根据条件确定执行哪一个分支。当有多个条件满足时,满足条件的分支并行执行。
name: test
steps:
- id: step1
branches:
- id: b1
# 分支条件表达式
if: '${{ a }} > 0'
steps:
- id: step2
- id: b2
# 默认分支
# 当其他条件全为 false 时执行
steps:
- id: step3
- id: step4
分支包含的属性有:
| key | 名 称 | 说 明 |
|---|---|---|
| id | 标识 | 节点的唯一标识 |
| name | 名称 | 有意义的名称,可以是中文等任意字符 |
| tag | 标签 | 标签设置 |
| if | 条件 | 分支条件表达式,使用 ${{ var }} 语法引用变量 |
| needs | 前置 | 前置分支id列表,只有前置分支完成后才执行当前分支,可以通过 needs 将分支执行线性化 |
| vars | 变量 | 本地变量定义 |
| steps | 步骤列表 | 分支内的步骤列表 |
| inputs | 输入 | 输入schema定义 |
| outputs | 输出 | 输出schema定义 |
分支前置依赖 (needs)
当分支设置了 needs 后,该分支会进入 Pending 状态,等待前置分支完成后再执行:
branches:
- id: b1
if: 'true'
name: 分支1
steps:
- id: step11
uses: acts.core.irq
params:
key: act1
- id: b2
if: 'true'
name: 分支2
needs:
- b1
steps:
- id: step21
活动
活动(Act)是最小的执行单元,代表一个具体的操作。活动通过 uses 指定使用的包,通过 params 传递参数。
活动属性如下:
| key | 名 称 | 说 明 |
|---|---|---|
| id | 标识 | 活动的唯一标识 |
| name | 名称 | 有意义的名称,可以是中文等任意字符 |
| tag | 标签 | 标签设置 |
| rn | 资源名 | 用于权限控制的资源名称 |
| uses | 包名 | 使用的包名称 |
| params | 参数 | 传递给包的参数 |
| inputs | 输入 | 输入数据 |
| outputs | 输出 | 导出数据 |
| options | 选项 | 额外选项,如 exposes 导出变量 |
| metadata | 元数据 | UI元数据,不发送给客户端 |
活动类型
活动根据 uses 指定的包来决定其行为:
| 包名 | 类型 | 说明 |
|---|---|---|
acts.core.irq | 中断请求 | 由引擎发起中断,等待客户端响应后继续 |
acts.core.msg | 消息 | 发送消息到客户端,不需要响应 |
acts.transform.set | 设置 | 设置变量值 |
acts.transform.code | 代码 | 执行 JavaScript 代码 (QuickJS) |
acts.core.block | 块 | 包含子活动列表,按模式执行 |
acts.core.parallel | 并行 | 对集合进行并行执行 |
acts.core.sequence | 顺序 | 对集合进行顺序执行 |
acts.core.subflow | 子流程 | 调用另一个工作流 |
acts.core.action | 命令 | 执行引擎命令 (error, complete 等) |
活动示例
# 中断请求 - 等待客户端完成
- uses: acts.core.irq
params:
key: act1
# 消息 - 单向通知客户端
- uses: acts.core.msg
params:
key: msg_notify
# 设置变量
- uses: acts.transform.set
params:
a: 10
b: hello
# 执行脚本
- uses: acts.transform.code
params: |
let x = $get("a");
$set("result", x * 2);
Set 设置变量
活动 acts.transform.set 用来设置当前变量值。
name: test
steps:
- id: step1
uses: acts.transform.set
params:
a: 5
b: hello
如果当前父节点或全局有相同名称的变量,则更新该变量的值:
name: test
vars:
- name: a
value: 0
steps:
- id: step1
# 此活动将全局变量 'a' 更新为 5
uses: acts.transform.set
params:
a: 5
IRQ 中断请求
活动 acts.core.irq 由引擎发起中断请求,等待客户端响应。当活动发起时,活动处于中断状态(interrupted),需要客户端接受并处理后调用完成接口继续执行。
name: test
steps:
- id: step1
uses: acts.core.irq
params:
# 活动标识key,发送给客户端的消息key为此值
key: act1
options:
# 导出变量到父节点
exposes:
- name: v
活动属性如下:
| key | 名 称 | 说 明 |
|---|---|---|
| id | 标识 | 活动节点的唯一标识 |
| name | 名称 | 有意义的名称,可以是中文等任意字符 |
| tag | 标签 | 标签设置 |
| uses | 包名 | 固定为 acts.core.irq |
| params | 参数 | 传递给包的参数,主要包含 key 作为消息关键字 |
| inputs | 输入 | 发送给客户端的输入数据 |
| outputs | 输出 | 客户端返回后导出到父节点的变量 |
| options | 选项 | 额外选项,如 exposes 导出变量 |
客户端处理
客户端通过 on_message 回调接收中断请求消息,处理后调用 do_action2 完成:
#![allow(unused)]
fn main() {
engine.channel().on_message(move |e| {
if e.is_params_key("act1") && e.is_state(MessageState::Created) {
// 处理业务逻辑
// 完成后调用 Next 继续执行
rt.do_action2(&e.pid, &e.tid, EventAction::Next, Vars::new()).unwrap();
}
});
}
MSG 消息
活动 acts.core.msg 由引擎发起单向消息,发送给客户端后不等待响应,立即继续执行后续步骤。
name: test
steps:
- id: step1
uses: acts.core.msg
params:
# 消息关键字key
key: msg1
inputs:
a: 1
活动属性如下:
| key | 名 称 | 说 明 |
|---|---|---|
| id | 标识 | 活动节点的唯一标识 |
| name | 名称 | 有意义的名称,可以是中文等任意字符 |
| tag | 标签 | 标签设置 |
| uses | 包名 | 固定为 acts.core.msg |
| params | 参数 | 传递给包的参数,主要包含 key 作为消息关键字 |
| inputs | 输入 | 发送给客户端的输入数据 |
客户端接收
客户端通过 on_message 接收消息:
#![allow(unused)]
fn main() {
engine.channel().on_message(move |e| {
if e.is_params_key("msg1") && e.is_state(MessageState::Completed) {
// 消息已收到,直接处理即可,无需响应
println!("收到消息: {:?}", e.inputs);
}
});
}
Block 块
使用 acts.core.block 将多个活动组合成一个块执行。支持 sequence(顺序)和 parallel(并行)两种模式。
name: test
steps:
- id: step1
uses: acts.core.block
params:
# 执行模式: sequence (顺序) 或 parallel (并行)
mode: sequence
acts:
- uses: acts.transform.set
params:
count: 0
- uses: acts.core.irq
params:
key: act1
- uses: acts.core.msg
params:
key: done
模式对比
| 模式 | 说明 |
|---|---|
sequence | 按顺序逐个执行子活动,后一个等待前一个完成 |
parallel | 所有子活动同时并行执行 |
变量导出
块内活动可以通过 options.exposes 将变量导出到父节点:
steps:
- id: step1
uses: acts.core.block
params:
mode: sequence
acts:
- uses: acts.core.irq
params:
key: act1
options:
exposes:
- name: result
Parallel 并行执行
使用 acts.core.parallel 对集合进行并行执行,所有子活动同时发起,互不依赖。
name: test
steps:
- id: step1
vars:
- name: items
value:
- u1
- u2
- u3
uses: acts.core.parallel
params:
in: '${{ items }}'
acts:
# 会生成 3 个 irq 活动,同时并行执行
- uses: acts.core.irq
params:
key: act1
与 sequence、block 的区别
| 类型 | 包 | 说明 |
|---|---|---|
| 并行 | acts.core.parallel | 所有子活动同时并行执行 |
| 顺序 | acts.core.sequence | 子活动按顺序逐个执行,后一个等待前一个完成 |
| 块 | acts.core.block | 按 mode: sequence 或 mode: parallel 模式执行嵌套 acts |
变量注入
引擎会自动将 index 和 value 注入到每个子活动的变量上下文中,可以在子活动中通过 ${{ index }} 和 ${{ value }} 访问。
代码生成集合
可以结合 acts.transform.code 动态生成集合:
steps:
- id: step1
uses: acts.transform.code
params: |
let list = ["u1", "u2", "u3"];
$set("items", list);
- id: step2
uses: acts.core.parallel
params:
in: '${{ items }}'
acts:
- uses: acts.core.irq
params:
key: act2
Sequence 顺序执行
使用 acts.core.sequence 对集合进行顺序链式执行,后一个序列的执行依赖上一个序列的完成。
name: test
steps:
- id: step1
vars:
- name: items
value:
- u1
- u2
uses: acts.core.sequence
params:
in: '${{ items }}'
acts:
# 会生成 2 个 irq 活动,按顺序逐一执行
- uses: acts.core.irq
params:
key: act1
与 parallel 的区别
| 类型 | 包 | 说明 |
|---|---|---|
| 并行 | acts.core.parallel | 所有子活动同时并行执行 |
| 顺序 | acts.core.sequence | 子活动按顺序逐个执行,后一个等待前一个完成 |
| 块 | acts.core.block | 按 mode: sequence 模式顺序执行嵌套 acts |
引擎会自动将 index 和 value 注入到每个子活动的变量上下文中。
Subflow 子流程
使用 acts.core.subflow 调用另一个工作流模型(子流程)。
name: test
steps:
- id: step1
uses: acts.core.subflow
params:
# 子流程的模型 ID
to: sub_workflow_id
# 传递给子流程的输入数据
a: '${{ value }}'
子流程定义
子流程是一个独立的工作流模型:
id: sub_workflow_id
name: sub_flow
inputs:
type: object
properties:
a:
type: integer
outputs:
type: object
properties:
result:
type: string
steps:
- id: sub_step1
uses: acts.core.irq
params:
key: sub_act
- id: sub_step2
uses: acts.core.msg
params:
key: sub_done
数据传递
子流程的输入通过 params 传递,子流程的输出通过 options.exposes 导出回父流程:
steps:
- id: step1
uses: acts.core.subflow
params:
to: sub_workflow_id
input_value: '${{ parent_var }}'
options:
exposes:
- name: result
Action 命令
使用 acts.core.action 执行引擎命令,如触发错误、完成步骤等。
steps:
- id: step1
uses: acts.core.irq
params:
key: act1
timeouts:
# 超时后触发错误
- uses: acts.core.action
if: $cost_in('8s')
params:
action: error
options:
ecode: err_timeout
支持的命令
| 命令 | 说明 |
|---|---|
error | 触发错误,可传递 ecode 指定错误码 |
客户端命令
客户端也可以通过 do_action2 执行以下操作来影响活动状态:
| 操作 | EventAction | 说明 |
|---|---|---|
| 完成 | Next | 完成当前活动,继续下一步 |
| 提交 | Submit | 提交当前活动 |
| 回退 | Back | 回退到指定步骤 |
| 取消 | Cancel | 取消指定活动 |
| 跳过 | Skip | 跳过当前活动 |
| 中止 | Abort | 中止当前活动 |
| 错误 | Error | 标记活动为错误 |
| 移除 | Remove | 移除活动 |
#![allow(unused)]
fn main() {
// 完成活动
rt.do_action2(&pid, &tid, EventAction::Next, Vars::new()).unwrap();
// 触发错误
let mut options = Vars::new();
options.set("ecode", "err1");
rt.do_action2(&pid, &tid, EventAction::Error, options).unwrap();
// 回退到指定步骤
let mut options = Vars::new();
options.set("to", "step1");
rt.do_action2(&pid, &tid, EventAction::Back, options).unwrap();
}
Code 代码执行
使用 acts.transform.code 执行 JavaScript 代码(QuickJS 引擎),可以在流程中进行变量计算、数据转换、条件判断等。
steps:
- id: step1
uses: acts.transform.code
params: |
let x = $get("a");
let y = $get("b");
$set("sum", x + y);
$set("message", "计算结果: " + (x + y));
内置函数
| 函数 | 说明 |
|---|---|
$get("key") | 获取变量值 |
$set("key", value) | 设置变量值 |
$ecode() | 获取当前错误码 |
$cost_in('2s') | 判断时间是否超过指定值 |
$inputs() | 获取上一步的输入数据 |
$data() | 获取当前数据 |
$env("key") | 获取环境变量 |
使用场景
变量计算:
- uses: acts.transform.code
params: |
let count = $get("count") || 0;
$set("count", count + 1);
数组操作:
- uses: acts.transform.code
params: |
let a = ["u1", "u2"];
let b = ["u2", "u3"];
$set("merged", a.concat(b));
条件判断与错误:
- uses: acts.transform.code
params: |
if ($get("status") != "ok") {
$set("ecode", "invalid_status");
}
案例
以下示例展示了 acts 工作流引擎的各种使用场景。
基础示例
| 示例 | 说明 | 路径 |
|---|---|---|
| 简单循环 | 使用 JavaScript 代码实现循环累加 | examples/simple |
| While 循环 | 使用 while 条件实现循环累加 | examples/while |
| 程序构建 | 使用 Rust Builder API 构建工作流 | examples/model_build |
交互示例
| 示例 | 说明 | 路径 |
|---|---|---|
| 动作交互 | 使用 IRQ 中断请求与客户端交互 | examples/actions |
| 审批流程 | 多角色审批工作流(PM、GM) | examples/approve |
| 消息通知 | 使用 MSG 消息发送单向通知 | examples/message |
错误与超时
| 示例 | 说明 | 路径 |
|---|---|---|
| 异常处理 | 使用 catches 捕获和处理错误 | examples/catches |
| 超时处理 | 使用 timeouts 处理步骤超时 | examples/timeout |
高级特性
| 示例 | 说明 | 路径 |
|---|---|---|
| 事件驱动 | 使用 on 事件触发工作流启动 | examples/event |
| 子流程 | 使用 subflow 调用子工作流 | examples/subflow |
| 自定义包 | 创建和注册自定义包 | examples/package |
| 自定义变量 | 注册和使用自定义用户变量 | examples/user_var |
插件示例
| 示例 | 说明 | 路径 |
|---|---|---|
| HTTP 请求 | 使用 acts-package-http 发送 HTTP 请求 | examples/plugins/http |
| Shell 执行 | 使用 acts-package-shell 执行 Shell 脚本 | examples/plugins/shell |
| 状态管理 | 使用 acts-package-state 管理状态 | examples/plugins/state |
运行示例
# 运行审批流程示例
cargo run --example approve
# 运行异常处理示例
cargo run --example catches
# 运行超时处理示例
cargo run --example timeout