Rivet Workflow Engine 状态模型全解析:从状态机到持久化与读取
Rivet Workflow Engine 状态模型全解析从状态机到持久化与读取【免费下载链接】actorsRivet Actors are the primitive for stateful workloads. Built for AI agents, collaborative apps, and durable execution.项目地址: https://gitcode.com/GitHub_Trending/riv/actorsWorkflow Engine 是 Rivet 生态中面向 TypeScript 的持久化执行引擎其核心设计是状态即数据工作流的每一个运行瞬间都通过引擎驱动EngineDriver持久化从而在进程重启、崩溃或部署迁移后能够确定性恢复。本文基于 state.md 展开结合引擎源码与测试系统讲解工作流的七种生命周期状态、持久化存储的内容与二进制编码、通过 handle 读取状态的方法以及如何编写真正可恢复的持久化工作流数据。读完本文你将能够准确判断工作流处于哪个阶段、知道什么数据被存到了哪里、以及如何安全地在步骤与循环之间传递可恢复的数据。工作流状态机七个生命周期状态Workflow Engine 将工作流的生命周期建模为一个显式的状态机。文档定义了七种状态在源码中以联合类型的形式被完整声明// rivetkit-typescript/packages/workflow-engine/src/types.ts export type WorkflowState | pending // 尚未开始 | running // 正在执行 | sleeping // 等待 deadline 或消息 | failed // 回滚后失败 | completed // 成功结束 | cancelled // 永久停止 | rolling_back; // 正在执行回滚处理器这七种状态覆盖了一条工作流从诞生到终局的全过程pending工作流实例已创建但尚未开始执行。在 storage.ts 的createStorage()中新存储的初始状态即为pendinghandle.getState()在驱动中读不到状态键时也会返回pending见 index.ts。running正在向前执行。executeWorkflow在调用用户工作流函数前会显式写入storage.state running见 index.ts。sleeping等待一个截止时间sleep、外部消息queue wait或重试退避。此时工作流会flush状态并调用driver.setAlarm(workflowId, deadline)把唤醒责任交给调度器见 index.ts。rolling_back工作流遇到不可恢复错误后进入回滚阶段逐个逆序执行已注册的 rollback 处理器见 index.ts。failed回滚完成后最终失败错误元数据被持久化见 index.ts。completed工作流函数正常返回输出被持久化同时清除待定的闹钟driver.clearAlarm见 index.ts。cancelled通过handle.cancel()永久停止写入cancelled状态并清除所有闹钟见 index.ts。从源码结构看状态的推进发生在 executeWorkflow 的收尾逻辑中成功路径写completedSleepError转sleepingStepFailedError重试退避也转sleeping并设置重试闹钟不可恢复错误先置rolling_back执行回滚回滚结束后置failed。测试 storage.test.ts 验证了成功与失败两条路径的最终落盘状态。引擎持久化了哪些数据工作流状态并非只包含一个状态枚举。引擎围绕每个工作流实例持久化以下几类数据工作流输入Workflow input仅在首次运行时捕获。恢复执行时优先读取已存储的输入以保证确定性重放——这是同样的输入产生同样的执行路径的根基。对应WORKFLOW_FIELD.INPUT键见 keys.ts 与 index.ts。工作流输出Workflow output工作流正常完成时持久化可通过handle.getOutput()读取。错误元数据Workflow error metadata失败时记录WorkflowError含 name、message、stack、自定义 metadata对应 types.ts 中的结构。历史条目History entries为 step、loop、sleep、join、race、message 以及 rollback checkpoint 各自记录持久化条目。这些条目共同构成了重放所需的完整执行轨迹。从数据组织上看Storage接口types.ts把上述内容划分为名称注册表nameRegistry、历史history、条目元数据entryMetadata、输出output、状态state与错误error并为每个持久化字段保留了flushed*影子值用于增量落盘。存储层的二进制编码与键布局所有状态数据都通过驱动EngineDriver以 KV 形式落盘。为了让历史按确定性顺序重放引擎使用fdb-tuple对键进行元组编码并用整数前缀区分不同数据类别见 keys.ts前缀数据类别键结构1名称注册表[1, index]2历史条目[2, ...locationSegments]3工作流元数据[3, field]field 1state、2output、3error、4input4条目元数据[4, entryId]驱动接口driver.ts只要求实现get/set/delete/batch/list与setAlarm/clearAlarm等原语其中list返回结果必须按键的字典序排序否则会导致名称注册表重建与历史重放不确定。每个工作流实例运行在独立的 KV 命名空间中引擎是执行期间该命名空间唯一的读写者——隔离由宿主系统如 Cloudflare Durable Objects 或专用 actor 进程提供。加载与落盘分别在 storage.ts 的loadStorageL140-L196与flushL240-L357中实现flush只写入新增名称、dirty 条目、dirty 元数据以及发生变化的状态/输出/错误并通过driver.batch一次性提交对于声明了atomicBatch的驱动整个批次会被视为一个不可分割的单元。批次大小有硬性约束MAX_KV_BATCH_ENTRIES 128、MAX_KV_BATCH_PAYLOAD_BYTES 976KB普通驱动会被自动拆分成多个合规批次而 storage.test.ts 验证了拆分逻辑与原子驱动的行为差异。条目及其元数据的二进制编码定义在 v1.bareBARE 协议模式用户数据统一以 CBOR 编码的任意二进制 blob 存储状态枚举、条目类型、join/race 分支状态与回滚检查点都有明确的 schema 定义。通过 handle 读取当前状态与输出在运行侧用户通过runWorkflow()返回的 handle 查询工作流状态const state await handle.getState(); // WorkflowState const output await handle.getOutput(); // TOutput | undefined对应的实现非常直接index.tsgetState()从驱动读取状态键[3, 1]getOutput()读取输出键[3, 2]并反序列化。完整的工作流句柄接口定义在 types.ts 中除了读取状态与输出还提供handle.result等待工作流完成或让出yield的 Promise解析后携带{ state, output, sleepUntil?, waitingForMessages? }handle.message(name, data)向等待消息的工作流发送消息handle.wake()立即唤醒live 模式下直接调用内存中的 sleep waiteryield 模式下设置当前时间闹钟handle.recover()把失败/耗尽的条目元数据重置为 pending 并清除错误重新调度执行handle.evict()请求优雅退出可被其他节点恢复handle.cancel()永久取消写入cancelled并清除闹钟。测试 handle.test.ts 在 yield 与 live 两种模式下验证了完成后的getState()返回completed、getOutput()返回工作流返回值。持久化工作流数据没有共享可变状态文档强调一个关键约束工作流内部不存在可变的共享状态。模块级变量、全局变量在重启后必然丢失因为它们从不被持久化。要构造可恢复的数据必须依赖引擎提供的两条路径通过ctx.step()的返回值传递步骤的执行结果被持久化为历史条目重放时直接从历史中恢复而不会重复执行副作用代码避免重启后重复扣款这类问题。使用ctx.loop()的持久化迭代状态循环维护一个跨迭代的state字段每轮迭代把新状态写入LoopEntry.state重启后从已保存的状态继续而非从头开始。从源码看LoopEntry正是为此设计的types.ts{ state, iteration, output? }三个字段全部持久化。循环还支持historyPruneInterval默认 20周期性裁剪旧迭代历史控制存储体积。配合 architecture.md 中描述的 location 系统与 NameIndex 优化每个 step/loop 的名字只存储一次历史条目通过位置路径索引保证重放时能精确对应到代码中的调用点。状态读取的实战建议在实际项目中使用handle.getState()/handle.getOutput()时有几条与状态模型直接相关的注意点输出只在完成后存在getOutput()在未完成时返回undefined。如果下游逻辑需要区分尚未完成与完成了但输出为 undefined应结合getState()或handle.result一起判断。sleeping 不代表失败sleeping是等待闹钟或消息的正常中间态外部系统可以通过handle.wake()提前唤醒。failed 后可以 recoverrecover()会把工作流状态从failed重置为sleeping并重新调度适合人工介入修复后重试的场景。确定性优先工作流函数在 step 之外必须保持确定性不要用Math.random()、Date.now()所有 I/O 与副作用放进 step否则重放时的状态可能与首次执行分叉触发HistoryDivergedError。更多关于驱动要求、存储 schema、loop 状态管理与历史裁剪的实现细节可继续阅读 architecture.md完整的七种状态声明与 handle 方法说明见 QUICKSTART.md状态相关的持久化测试见 storage.test.ts 与 handle.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),仅供参考