让工作流持久化
保存检查点、在不占用进程的情况下等待、恢复租约,并按照显式策略继续运行。
持久化模型
持久工作流会保存有语义的进度,再由任意兼容 Worker 继续执行。进程是否持续在线不再 影响正确性。
需要跨发布存活、等待人员或定时器、协调昂贵副作用,或在 Worker 消失后恢复的任务, 都适合使用持久化。
在本机运行一个持久 Task
单机持久 Worker 启用 SQLite;如果多个进程或主机需要并发 Claim,改用 PostgreSQL。
[dependencies]
runifold = { version = "=0.9.0", features = ["sqlite-bundled"] }
serde_json = "1"
tokio = { version = "1", features = ["macros", "rt-multi-thread"] }假设 workflow 是在 Workflow 指南中构建的 Workflow,下面代码会
创建持久控制面、注册准确版本、入队一个 Task,并运行一次 Claim:
use std::{sync::Arc, time::Duration};
use runifold::{
Budget, CapabilitySet, LeaseDuration, WorkerId, WorkflowDefinition,
WorkflowRegistry, WorkflowStore, WorkflowTask, WorkflowWorker,
sqlite::SqliteWorkflowStore,
};
use serde_json::json;
let store = Arc::new(SqliteWorkflowStore::open("runifold-workflows.db")?);
let mut registry = WorkflowRegistry::new();
registry.register(WorkflowDefinition::new(
Arc::new(workflow),
Budget::default(),
CapabilitySet::new(),
))?;
store
.enqueue(WorkflowTask::new(
"order-intake",
1,
json!({ "order_id": "ord_42" }),
)?)
.await?;
let worker = WorkflowWorker::new(
store,
registry,
WorkerId::parse("worker-local-1")?,
LeaseDuration::new(Duration::from_secs(30))?,
Duration::from_secs(10),
)?;
let outcome = worker.run_once().await?;
println!("{outcome:?}");run_once() 会返回 Idle、Completed、Retried、Suspended、Failed、
DefinitionUnavailable 或 LeaseLost。生产服务通常使用有界 Supervisor;测试或接入
现有 Job Loop 时可以直接使用 run_once()。
检查点与版本
Checkpoint 记录工作流定义、当前阶段、已完成步骤输出、预算状态和恢复元数据。 Revision 可防止两个 Worker 提交互不兼容的进度。
步骤输出应可序列化、紧凑且有版本。大型产物放在 Checkpoint 外,通过不可变 ID 引用。
定时器、信号与人工审核
Wait 会持久化等待意图,不占用线程或进程。Timer 在截止时间后唤醒;Signal 在审批等 外部事件发生时唤醒。
Signal 名称与 Payload 是应用契约。验证发送者身份、按 Signal ID 去重,并保留足够 历史来解释谁在何时、为何恢复了工作流。
Worker 与租约
Worker 通过有限 Lease 认领任务。健康 Worker 会续租;租约过期后,其他 Worker 可以 恢复任务。
SQLite Store 适合本地持久执行;分布式 Worker 和多进程协调使用 PostgreSQL 工作流 Feature。
恢复策略
恢复必须区分可以安全重复的步骤,以及结果未知的 Effect。按边界配置 Resume Policy, 在继续前先对账不确定副作用。
旧任务仍存在时,新版本必须保持工作流定义兼容。把 Step ID、序列化状态与 Signal 名称当作数据库 Schema 管理。
恢复演练
上线前应对真实 Store 完成以下演练:
- 使用已知 Checkpoint ID 入队一个 Task;
- 在确定性 Step 执行中终止 Worker;
- 等待超过 Lease 时长,再用新的 Worker Identity 启动;
- 验证已完成 Step 不会重复,Usage 不会减少;
- 在外部写入边界重复演练,确认结果依据 Effect 证据被重放或标记为 Ambiguous, 而不是根据超时猜测;
- 在版本 1 Task 等待期间部署版本 2,确认两个 Definition 都仍被注册。
常见失败
| Outcome 或错误 | 含义 | 修复 |
|---|---|---|
DefinitionUnavailable | Worker 缺少 Task 指定的名称/版本 | 部署或重新注册该准确版本 |
LeaseLost | 另一个 Owner 可能已经接管 | 立即停止写入,让当前 Owner 恢复 |
持续 Retried | 失败策略或租户预算推迟任务 | 检查类型化失败和 Retry Delay |
| Checkpoint Conflict | 过期 Revision 尝试提交 | 丢弃旧状态并重新加载 |
| Ambiguous In-flight | 外部是否完成未知 | 先对账,禁止静默成功或重跑 |
| Task 一直 Waiting | Timer 未到期或 Signal 未接受 | 检查 Wake Time、Tenant、Signal 名与去重 ID |
对应控制面细节继续阅读 Worker、 Wait 与 Signal、Checkpoint 恢复和 Workflow 版本。