组合确定性工作流
连接类型化步骤、并行分支与首次成功竞速,无需让模型控制每一次跳转。
Agent 还是工作流
下一步确实需要模型判断时使用 Agent;转换规则可以确定表达时使用工作流。多数生产 系统会把两者组合起来:确定性编排包围有界的 Agent 步骤。
审批、补偿、并行展开与完成规则因此留在类型化代码中,而不是藏在 Prompt 里。
安装并运行最小 Workflow
Workflow 已包含在默认的 runtime Feature 中。先用一个不依赖模型的步骤理解编排,
再加入 Provider 和 Agent。
[dependencies]
anyhow = "1"
futures-executor = "0.3"
runifold = "=0.9.0"
serde_json = "1"use anyhow::Context;
use runifold::{
Budget, BudgetTracker, CapabilitySet, RunContext, Workflow, WorkflowStep,
WorkflowStepError, WorkflowStepFuture,
};
use serde_json::{Value, json};
struct NormalizeOrder;
impl WorkflowStep for NormalizeOrder {
fn execute<'a>(
&'a self,
input: Value,
_run: &'a RunContext,
) -> WorkflowStepFuture<'a> {
Box::pin(async move {
let order_id = input["order_id"]
.as_str()
.ok_or_else(|| WorkflowStepError::Execution("missing order_id".to_owned()))?;
Ok(json!({ "order_id": order_id, "normalized": true }))
})
}
}
fn main() -> anyhow::Result<()> {
let workflow = Workflow::builder("order-intake")
.version(1)
.step("normalize", NormalizeOrder, CapabilitySet::new())
.build()
.context("failed to build workflow")?;
let run = RunContext::root(
BudgetTracker::new(Budget::default()),
CapabilitySet::new(),
);
let outcome = futures_executor::block_on(
workflow.run(json!({ "order_id": "ord_42" }), &run),
)?;
println!("{}", outcome.output);
Ok(())
}运行 cargo run。一个步骤的输出会成为下一个步骤的输入;outcome.output
是最终规范结果,outcome.usage 是实际消耗的预算。构建失败表示图定义无效;运行失败
则表示步骤、Capability、预算、Deadline 或取消边界阻止了执行。
顺序步骤
工作流定义为每一步提供稳定 StepId、类型化输入输出和明确失败策略。只把下一步
需要的数据传过去。
好的步骤边界可以单独重试和观察,例如读取记录、分类、申请审批、执行 Effect、 保存结果。不要用一个巨大步骤混合所有职责。
并行分支
独立任务可以并行以降低延迟。启动前预留预算,并在 Worker 或租户边界限制并发。
明确选择汇合规则:要求全部成功、接受部分结果,或返回类型化的部分成功。一条慢分支 不能静默延长整体 Deadline。
首个成功竞速
冗余读取或可替换模型路由适合首个成功竞速:第一个有效结果获胜,其他分支被取消。
不要竞速非幂等写操作。取消无法撤回已经到达外部目标的副作用。
加入 Agent 步骤
先构建 Agent,再用 Arc 包装,并只授予该步骤真正需要的 Capability。plan、
write 这样的稳定 Step ID 会成为持久身份;只要旧 Checkpoint 仍可能恢复,就不要
重命名它们。
let workflow = Workflow::builder("plan-and-write")
.version(1)
.agent("plan", planner, CapabilitySet::new())
.agent("write", writer, CapabilitySet::new())
.build()?;
let outcome = workflow.run("设计一个重试策略", &run).await?;
println!("{}", outcome.output);确定性的 Rust 逻辑使用普通 .step(...);只有真正需要模型判断时才用
.agent(...)。校验、写入、审批和状态迁移应留在 Prompt 之外。
加入审查门控生成
当应用自有 Reviewer 必须在 Workflow Commit 前批准生成值时,使用
.repairable_agent(...):
let workflow = Workflow::builder("reviewed-answer")
.repairable_agent(
"draft",
answer_agent,
compliance_reviewer,
WorkflowRemediationPolicy::new(2),
generation_capabilities,
reviewer_capabilities,
)
.build()?;第一次生成接收普通 Workflow 输入。Repair Verdict 会先持久化被拒绝的 Candidate 与 结构化 Feedback,再进入下一次生成。Generation、Review、Approval 与 Repair 是独立的 Write-ahead 子阶段,因此恢复会从持久 Candidate 继续,不会再次生成。升级时必须先排空 0.7 Worker:0.8/0.9 写入的 Checkpoint Schema v5 无法被 0.7 Worker 读取。
选择执行模式
| 需求 | Builder 形式 | 关键约束 |
|---|---|---|
| 每个阶段依赖前一步结果 | 连续 step / agent | 每个输出都应可序列化 |
| 必须获得所有独立结果 | parallel | 为每个分支预留有限预算 |
| 第一个有效只读结果即可 | race | 只允许 Pure 或 ReadOnly Capability |
| 进程退出后仍要继续 | 持久 Store + Worker | 固定定义版本与稳定 ID |
| 等待人工或定时恢复 | 持久 Wait / Signal | 认证并去重 Signal |
排查失败的 Workflow
- 在启动时构建 Workflow,遇到重复或无效 ID 立即失败。
- 记录 Workflow 名称、定义版本、Step ID、Run ID 与类型化错误种类。
- 先检查
outcome.usage和 Run Journal,不要直接认定是 Provider 故障。 - 如果外部写入可能已经完成,先对账,不要盲目重跑。
- 对持久执行,确认已部署 Worker 仍注册了 Checkpoint 所指定的准确版本。
需要跨重启执行时继续阅读持久 Workflow,扇出与 Race 阅读并行执行,队列与租约运行阅读 Worker。