协调 Timer、Signal 与持久等待
让工作流在不占用进程的情况下暂停,幂等投递 Signal,治理长期等待并安全恢复。
等待为何是持久状态
长期业务流程的大部分时间都在等待:送达日期、审批、Webhook、人工回答或外部 Batch。 让 Worker 或 Async Task 一直存活会浪费容量,并在重启时丢失状态。
持久 Wait 提交暂停原因和允许唤醒它的事件,然后 Worker 释放所有权,之后可由另一进程 恢复。
Timer
Timer 保存由 Store 决定的权威唤醒时间。在此之前 Claim Query 不会返回它;到期后无需 进程 Sleep 或内存 Scheduler,就会重新就绪。
Timer 适合重试计划、提醒、轮询间隔和业务期限。它与操作 Timeout 不同:Timer 决定何时 恢复,Timeout 限制一次尝试。
Signal
Signal 是发送给等待 Workflow 的外部输入。每次投递都要有稳定 Identity,使网络重试不会 产生重复跳转。Signal 应持久化,并按 Tenant 与 Checkpoint 隔离。
明确 Early、Duplicate、Unknown 与 Late Signal 的行为。不能把不透明 Task ID 或 Signal ID 本身当作授权。
Interrupt
Interrupt 把 Workflow 转成显式 input_required 状态,标识所需决策与接受的响应格式。
Approve、Edit、Reject 或 Cancel 命令必须幂等。
人工审核不是一种特殊 Prompt,而是具有 Identity、授权、保留与可审计结果的持久协议 边界。
保留与治理
限制 Wait 时长、Pending Signal 数、轮询频率与保留历史。Supervisor 应找出陈旧 Wait 并 暴露问题,而不是虚构完成。
协议 TTL、执行 Deadline、业务 Due Date 与物理删除 Retention 是不同概念。分别建模; 混用可能取消仍在运行的工作,或无限期保存敏感状态。
构建 Timer、Signal 与人工审核 Wait
use std::time::Duration;
let workflow = Workflow::builder("approval")
.version(1)
.step("prepare", PrepareProposal, CapabilitySet::new())
.wait_for_signal_or_timeout(
"wait_for_documents",
"documents_ready",
Duration::from_secs(24 * 60 * 60),
)
.interrupt("manager_review", "Approve, edit, or reject this proposal")
.timer("cooldown", Duration::from_secs(30))
.step("commit", CommitDecision, write_capabilities)
.build()?;Worker 在释放 Lease 前持久化 Wait。通过 WorkflowStore 发布 Signal,使用稳定的
WorkflowSignalId、已认证 Tenant、Checkpoint ID、Name 与 Payload;重复投递应视为
同一事件。Interrupt 要展示带 Key 的 Input Request,并提交一次幂等 Approve / Edit /
Reject 决策,绝不能把缺少 Response 推断为批准。