运行持久工作流 Worker
排队 Workflow Task、认领带栅栏 Lease、维持心跳、恢复 Checkpoint,并跨进程接续工作。
Worker 模型
持久 Workflow Definition 描述确定性步骤。Workflow Task 保存某个 Definition、Input、 Version、Tenant 与 Checkpoint Identity 的执行状态。Worker 是可替换进程,负责认领就绪 Task、执行有界工作,并通过 Store 提交状态。
不要只把正确性放在进程内 Future 中;进程可能在完成外部操作后、报告完成前停止。
Task 生命周期
Task 在 Queued、Leased、Waiting、Terminal 或 Cancelled 状态间变化。Store 拥有权威时间 和跳转。Worker 认领有界 Lease,根据注册的名称与版本重建 Workflow,从 Checkpoint 恢复,再提交进度、进入持久等待或记录终止结果。
Definition Identity 必须稳定。部署新代码不能把旧 Checkpoint 解释成另一张流程图。
运行有界 Supervisor
构造持久 Workflow 指南中的 WorkflowWorker 后,使用有界并发
和优雅取消来运行:
use std::{sync::Arc, time::Duration};
use runifold::{
CancellationToken, WorkflowSupervisor, WorkflowSupervisorConfig,
};
let config = WorkflowSupervisorConfig::new(8)?
.with_backoff(Duration::from_millis(25), Duration::from_secs(5))?;
let supervisor = WorkflowSupervisor::new(Arc::new(worker), config);
let shutdown = CancellationToken::new();
let signal = shutdown.clone();
tokio::spawn(async move {
if tokio::signal::ctrl_c().await.is_ok() {
signal.cancel();
}
});
let report = supervisor.run(&shutdown).await;
println!("{report:?}");Concurrency 只限制当前进程内同时进行的 Claim-and-execute Cycle,不能替代 Tenant Admission、数据库连接池限制、Model 并发或下游限流。关停时,Supervisor 停止调度新 Cycle,并等待已经开始的 Cycle 完成 Lease Protocol。
带栅栏 Lease
Lease 包含 Tenant、Task、Worker、过期时间与单调递增 Fencing Token。Heartbeat 延长有效 所有权;过期后另一 Worker 可以用更高 Token 认领。
所有改变状态的写入都验证 Ownership 与 Fencing。即使旧进程的网络调用稍后返回,它也 不能在接管后提交。分布式过期以数据库时间为准,而不是进程时钟。
Checkpoint 恢复
Checkpoint 保存语义进度,而不是序列化 Rust Future。恢复根据持久 Phase 与 Effect 证据 决定下一步。已完成 Effect 可以 Replay 记录结果;不确定的非幂等 Effect 必须关闭式失败, 交给应用协调。
并行 Branch 在开始前预留共享预算。First-Success Race 只在语义安全时取消失败者;取消 不会撤销已经完成的外部写操作。
Worker 运维
运维 Worker 时需要:
- 有界 Claim Batch 与并发;
- 明显短于 Lease Duration 的 Heartbeat Interval;
- 停止新 Claim 并保存已有进度的优雅关停;
- Queue Delay、Lease Loss、执行时间、Wait 与失败指标;
- 注册每个可恢复 Definition Version;
- 把 Store 健康与时钟行为纳入故障处理。
只有 Store 与下游系统能承受额外并发时才扩容。更多 Worker 无法修复 Hot Tenant、缺失 幂等性或无界 Workflow。
部署 Runbook
- 先部署 Schema 变更与所有仍活跃的 Workflow Definition;
- 启动新 Worker,每个进程使用唯一且稳定的
WorkerId; - 增加并发前确认 Claim、Heartbeat 与 Completion 正常;
- 优雅停止旧 Worker,并至少等待一个 Lease Window,再判断旧 Owner 已全部退出;
- 检查 Queue Age、Retry、Definition-unavailable 与 Lease Loss;
- 只有不存在引用旧版本的 Queued、Leased、Waiting 或 Recoverable Task 时才能删除定义。
Lease Loss 增多时,降低并发并检查数据库 Pool Wait、Transaction Latency、Heartbeat Timing 与 Runtime Stall。Worker Idle 但 Queue 增长时,检查 Tenant Policy、到期 Timer、 Definition Registration,以及 Task 是否因 Budget Admission 持续延期。