Rivet Kit Workflow Engine 长运行工作流完整指南:sleep 让出、循环状态检查点与驱逐恢复的底层原理
【免费下载链接】actorsRivet Actors are the primitive for stateful workloads. Built for AI agents, collaborative apps, and durable execution.项目地址: https://gitcode.com/GitHub_Trending/riv/actors
本文以 long-running-workflows.md 为骨架,结合@rivetkit/workflow-engine包(位于 rivetkit-typescript/packages/workflow-engine)的真实源码与测试,系统讲解长运行工作流的完整技术栈:工作流如何通过sleep/queue.next让出执行权、短睡眠与长睡眠如何分流、循环如何按迭代做状态检查点并裁剪历史、工作进程被驱逐(eviction)时如何优雅让位并异地恢复,以及EngineDriver调度接口的可靠性要求。读完本文,你将能够为 AI Agent、协作应用与持久化执行场景编写可暂停数小时甚至数天、跨进程重启无缝续跑的工作流,并理解其底层驱动协议。
长运行工作流的本质:durable history 与 driver scheduler
Rivet Kit Workflow Engine 是一个 TypeScript 的持久化执行(durable execution)引擎:工作流就是普通的 async 函数,但它的每一次操作都会被记录为持久化历史(durable history),因此可以在进程崩溃、工作进程被驱逐、部署滚动、主动 sleep 等任意时刻被打断,并在稍后从历史断点精确续跑,而不是从头重放。
长运行工作流能够"暂停、睡眠、跨进程重启恢复",正是由两个机制协同驱动的:
- durable history(持久化历史):每个操作(step、sleep、loop 迭代、队列消息等)都以 entry 的形式写入由
EngineDriver提供的隔离 KV 命名空间。再次运行时,工作流先从历史中加载已有 entry,命中则直接返回记录的结果,不再执行用户回调。参见 architecture.md 中 "Replay Execution" 一节与 QUICKSTART.md 的 "History Tracking / Replay / Deterministic Execution" 三大支柱。 - workflow driver scheduler(驱动调度器):
EngineDriver的setAlarm(workflowId, wakeAt)与clearAlarm(workflowId)负责把"到点该醒"的工作流重新调度回运行器。工作流进入睡眠后并不占用任何 worker 线程,而是由调度器在 deadline 到达(或外部消息送达)时再次触发runWorkflow。
也就是说,长运行 = 历史让你"记得住进度",调度器让你"到点被叫醒",两者缺一不可。
Yielding Execution:显式让出执行权
让出执行权(yielding)是长运行工作流最核心的编程动作。你不需要自己实现任何状态机,只需要在等待处调用两个 helper:
// 等待 5 分钟(在 workflow 内部) await ctx.sleep("wait-5-min", 5 * 60 * 1000); // 等待一条名为 "approval" 的队列消息 const [message] = await ctx.queue.next<string>("wait-approval", { names: ["approval"], });当工作流让出时,runWorkflow并不会阻塞等待,而是立即返回一个WorkflowResult,其state为"sleeping":
interface WorkflowResult<TOutput> { state: WorkflowState; // "sleeping" | "running" | "completed" | ... sleepUntil?: number; // 长 sleep 的唤醒 deadline waitingForMessages?: string[]; // 正在等待的消息名 }后续由两件事触发下一次 run:
- driver alarm:长 sleep 设置的定时闹钟到点;
- message wake-up:外部系统通过
handle.message(name, data)投递消息后,运行器被唤醒重新执行工作流。
源码证据在 index.ts:runWorkflow内部捕获到SleepError时调用setSleepState,它会先将存储状态置为"sleeping"并 flush,再调用driver.setAlarm(workflowId, deadline),最后返回{ state: "sleeping", sleepUntil, waitingForMessages };捕获到MessageWaitError时走setMessageWaitState,只返回waitingForMessages而不设置 alarm(等待消息唤醒)。工作流函数的写法保持不变——下次运行时会从历史中读到这条 sleep entry 已 "completed",直接跳过等待。
Short vs Long Sleeps:两种睡眠路径的分流
并非所有睡眠都需要惊动调度器。引擎根据driver.workerPollInterval(worker 轮询间隔,毫秒)这个阈值把 sleep 分成两类:
- 短睡眠(剩余时间 <
driver.workerPollInterval):直接在内存中等待。工作流进程保持存活,用定时器睡完剩余时间,不产生 alarm、不涉及调度器,因此开销极小。 - 长睡眠(剩余时间 ≥
driver.workerPollInterval):通过 driver 设置 alarm(setAlarm),然后抛出SleepError把控制权交还给调度器。worker 立即空闲,可以去执行其他工作流或缩容。
这正是"工作流可以暂停数小时甚至数天而不占用 worker 内存"的原理:一旦进入长睡眠,进程里只剩一条持久化的 alarm 记录和一个"sleeping"状态,任何内存资源都已释放。
源码实现
在 context.ts 的executeSleep中可以看到明确的分流逻辑:
const now = Date.now(); const remaining = deadline - now; if (remaining <= 0) { // 已过期:直接标记 completed 并 flush ... return; } // 短睡眠:在内存中等待(可被驱逐打断) if (remaining < this.driver.workerPollInterval) { await this.sleepOrEvict(remaining); this.checkEvicted(); ... return; } // 长睡眠:让出给调度器 throw new SleepError(deadline);sleepOrEvict(context.ts)会同时注册一个setTimeout完成定时器和 abort 监听器:定时器到点则正常 resolve,abortSignal触发则reject(new EvictedError()),且无论哪种结果都会在finally中清理 timer 与 abort 监听器,避免在长生命周期的 run signal 上留下悬挂监听器(内存泄漏)。
ctx.sleepUntil(name, timestampMs)与sleep(name, durationMs)同源:sleep内部先算deadline = Date.now() + durationMs,再委托给sleepUntil(context.ts)。
测试验证
sleep.test.ts 用yield/live两种运行模式交叉验证了这些行为:
"should complete short sleep in memory":workerPollInterval = 1000时,await ctx.sleep("short-sleep", 10)直接完成,且断言driver.getAlarm("wf-1")为undefined(没有设置闹钟);"should yield on long sleep"(yield 模式):长 sleep 后result.state === "sleeping"且sleepUntil定义在未来;"should schedule and clear alarms for long sleep":长 sleep 期间 alarm 存在,工作流完成后driver.getAlarm变为undefined(运行结束时调用driver.clearAlarm,见 index.ts);"should resume after sleep deadline":第一次运行返回"sleeping",等 deadline 过后再次runWorkflow,直接从历史续跑并返回"completed"。
Checkpointing Loop State:循环状态按迭代持久化
长时间运行的循环(比如分批消费消息队列、分批处理数据游标)是长运行工作流最常见的形态。ctx.loop()为此提供了按迭代持久化状态的能力,而不是把整个循环塞进一次 step:
const total = await ctx.loop({ name: "process-batches", state: { cursor: null, count: 0 }, // 初始状态 historyPruneInterval: 20, // 每 20 次迭代持久化一次并裁剪旧历史 run: async (ctx, state) => { const batch = await ctx.step("fetch", () => fetchBatch(state.cursor)); if (!batch.items.length) { return Loop.break(state.count); // 退出循环,返回最终值 } await ctx.step("process", () => processBatch(batch.items)); return Loop.continue({ cursor: batch.nextCursor, count: state.count + batch.items.length, }); }, });关键语义有两点:
- 每
historyPruneInterval次迭代持久化循环状态:每次迭代结束都会把state与iteration写入 loop entry(entry.kind.data.state/.iteration),并在达到historyPruneInterval整数倍时触发一次带裁剪的 flush。崩溃后重放时,从最后持久化的迭代状态继续,而不是从头再跑。 - 超过
historySize的旧迭代被裁剪:historySize默认等于historyPruneInterval(默认值 20,见 context.ts 的DEFAULT_LOOP_HISTORY_PRUNE_INTERVAL = 20)。这样回滚(rollback)只重放最后保留的若干次迭代,长时间运行的循环不会积累无界历史,存储占用与回放时间都保持有界。
你还可以把historySize设置得比historyPruneInterval大,例如"每 20 次迭代裁剪一次,但保留最近 100 次迭代",以换取更深的回滚能力。architecture.md 的 "History Size" 一节给出了具体的裁剪示例:在迭代 40 处裁剪时(historyPruneInterval=20, historySize=20),迭代 0-19 被删除,迭代 20-39 保留。
裁剪的底层实现
collectLoopPruning(context.ts)只在currentIteration > historySize时工作:它通过buildLoopIterationRange构造一个半开区间[fromIteration, keepFrom)(keepFrom = currentIteration - historySize),把所有落在这个区间内的迭代 entry 及其元数据一并标记删除,然后连同本轮状态写一起通过flushStorageWithDeletions原子落盘。实现中还维护lastPrunedUpTo游标,只删除"新过期"的迭代,避免每次从 0 开始重扫。
值得注意的工程细节:达到裁剪点时,flush 被**延迟(deferred)**到下一次迭代开始前执行(deferredFlush机制,见 context.ts),使状态写入与用户迭代代码并行推进,减少 IO 停顿。对应的循环与裁剪行为在 loops.test.ts 中有一系列测试覆盖。
Handling Eviction:优雅处理工作进程被驱逐
在 Serverless / Actor 场景中,worker 可能因水平扩缩容或滚动部署而在任意时刻被回收。引擎把这种优雅回收称为eviction:不是杀掉工作流,而是请求它在安全点保存状态、交还控制权,然后由调度器在别的 worker 上续跑。
工作流内有两种方式来感知驱逐并安全收手:
ctx.abortSignal:传给支持AbortSignal的 API(如fetch(url, { signal: ctx.abortSignal })),由引擎统一触发;ctx.isEvicted():轮询检查是否已被驱逐。
官方推荐的长任务模式是"分块干活 + 每块检查驱逐":
await ctx.step("long-task", async () => { while (!ctx.isEvicted()) { await doChunkOfWork(); // 每次只做一小块工作 } });这样,在驱逐信号到来时,当前块完成后立即退出 step,避免把工作流卡死在无法中断的同步长任务上。
源码证据
isEvicted()的实现就是一行:return this.abortSignal.aborted;(context.ts);evict()通过this.abortController.abort(new EvictedError())触发(context.ts);handle.evict()则直接调用上下文链路上的同一 abort(index.ts);- 工作流内任何
await若因 abort 抛EvictedError,runWorkflow会捕获并走setEvictedState(index.ts):只做一次 flush 保存当前全部脏状态,然后返回{ state: storage.state },把调度权交还 scheduler。
驱逐的语义是"保存状态、安全让位、异地恢复"——与永久性的handle.cancel()(写"cancelled"状态并清除 alarm)有本质区别。相关行为由 eviction-cancel.test.ts 覆盖。
Driver Considerations:EngineDriver 调度接口的可靠性要求
长运行工作流对宿主系统(Host System)暴露的接口就是EngineDriver(driver.ts)。除了 KV 读写(get/set/delete/list/batch等),与"长时间运行"直接相关的是两个调度方法与一个阈值:
export interface EngineDriver { // ...KV 操作... // 设置闹钟:在 wakeAt 唤醒指定工作流 setAlarm(workflowId: string, wakeAt: number): Promise<void>; // 清除工作流上任何待触发的闹钟 clearAlarm(workflowId: string): Promise<void>; // worker 轮询间隔(毫秒):决定短/长睡眠的阈值 readonly workerPollInterval: number; // 消息驱动(queue.next / handle.message 依赖) readonly messageDriver: WorkflowMessageDriver; // live 模式下等待指定消息名到达 waitForMessages(messageNames: string[], abortSignal: AbortSignal): Promise<void>; }对长运行工作流而言,driver 实现必须满足以下可靠性要求:
- alarm 必须持久化可靠:
setAlarm写入的闹钟不能因调度器重启而丢失。长 sleep 可能横跨数小时甚至数天,期间调度进程可能多次重启,闹钟必须能从持久化存储中恢复并继续生效。 - 到期的闹钟必须归还给 runner:调度器到期触发时,应把该
workflowId作为"可运行任务"交回给 worker,由 worker 再次调用runWorkflow。这是"睡醒续跑"闭环的关键一步。 - 完成或取消时清除闹钟:工作流完成(index.ts)与
handle.cancel()(index.ts)都会调用driver.clearAlarm,driver 必须保证不再触发已结束的工作流。 workerPollInterval的取值会直接影响调度压力:值越小,越多的 sleep 走 alarm 路径(调度开销大、worker 更空闲);值越大,越多的 sleep 在内存等待(worker 占用时间长)。应按业务实际 sleep 分布权衡。list()必须按字典序返回:工作流引擎依赖 key 的有序性做确定性重放与名称注册表重建(见 architecture.md 的 "Driver Requirements" 一节),无序会导致非确定性重放,这也是长运行稳定性的隐性前提。
此外要注意引擎的隔离模型:每个工作流实例拥有完全独立的 KV 命名空间,引擎执行期间是唯一的读写者;外部系统只能通过WorkflowHandle.message()走消息驱动投递消息,不能在 KV 层直接改动工作流状态。宿主系统(如 Cloudflare Durable Objects、独立 Actor 进程)负责提供这个隔离边界。
让长运行工作流稳定运行的实践要点
综合原文档、QUICKSTART 与源码,总结如下实践清单:
- 等待一律用
ctx.sleep/ctx.queue.next,不要用原生setTimeout:只有经过引擎的操作才会进入历史、才能跨重启恢复。 - 长循环用
ctx.loop并按迭代持久化:配合historyPruneInterval/historySize,把历史与回放时间保持有界;不要用while (true)原生循环。 - 长任务要响应驱逐:在 step 内部分块执行并轮询
ctx.isEvicted(),或把ctx.abortSignal传给可取消的 IO;让 eviction 在毫秒级生效。 - 避免同步阻塞与不确定代码:工作流函数主体应保持确定性,非确定性/副作用(
Math.random()、Date.now()、外部 IO)都放进 step 内部,否则重放会产生历史分歧(HistoryDivergedError)。 - driver 把 alarm 当一等公民:持久化、可靠触发、到期归还、完成清除,四件事缺一不可,这是"暂停数小时甚至数天"能否兑现的底层保障。
总结
长运行工作流是 Rivet Kit Workflow Engine 面向 AI Agent、协作应用与持久化执行场景的核心能力:ctx.sleep/ctx.queue.next负责让出执行权,workerPollInterval阈值把短睡眠留在内存、长睡眠交给 driver alarm,ctx.loop按迭代做状态检查点并裁剪有界历史,ctx.isEvicted()/ctx.abortSignal让工作流在扩缩容与部署中优雅让位并异地恢复,而这一切都建立在EngineDriver可靠持久化 alarm 与隔离 KV 之上。想进一步深入,可以继续阅读 QUICKSTART.md(含完整 API 与示例)、architecture.md(存储 schema、Location 系统、消息投递模型)以及 sleep.test.ts、loops.test.ts、eviction-cancel.test.ts 等测试用例。
【免费下载链接】actorsRivet Actors are the primitive for stateful workloads. Built for AI agents, collaborative apps, and durable execution.项目地址: https://gitcode.com/GitHub_Trending/riv/actors
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考