Rivet Actors Rust SDK(rivetkit)设计约束与开发指南
【免费下载链接】actorsRivet Actors are the primitive for stateful workloads. Built for AI agents, collaborative apps, and durable execution.项目地址: https://gitcode.com/GitHub_Trending/riv/actors
Rivet Actors 将 AI Agent、协作文档、实时聊天等有状态工作负载抽象为可休眠、可持久化的 Actor 原语。本文聚焦开源仓库中 rivetkit-rust/packages/rivetkit/CLAUDE.md 所定义的 Rust SDK 设计约束,结合 rivetkit-rust/packages/rivetkit 的源码、测试与官方示例,系统讲解rivetkit的架构定位、trait 化 Actor 编程模型、状态持久化规则、事件循环兼容层以及职责边界,帮助 Rust 开发者正确、高效地使用该 SDK 编写有状态 Actor 应用。
一、rivetkit 的架构定位:rivetkit-core之上的薄类型封装
rivetkit不是一套独立的 Actor 运行时,而是对底层rivetkit-core的薄类型封装(thin typed wrapper)。这一点在包描述与源码中均有明确体现:
- Cargo.toml 将包描述为 "Rust SDK for RivetKit actors, actions, events, queues, and test harnesses",其核心依赖为
rivetkit-core、rivetkit-client,序列化采用ciborium(CBOR),异步运行时为tokio。 - src/lib.rs 大量
pub use rivetkit_core::...直接转出底层类型,说明封装层的核心职责是类型化与边界编解码,而非重新实现运行时逻辑。
该封装策略对应了 rivetkit-rust/CLAUDE.md 中定义的 RivetKit 运行时边界:跨运行时的字节边界统一使用Vec<u8>形状的数据,SQL 边界类型显式共享,避免从 NAPI 专属的数据库包装器推导运行时 API 契约。rivetkit与rivetkit-typescript保持尽力而为的 API 对齐(best-effort parity),使得同一套 Actor 概念在两个语言 SDK 中可以映射到等价的结构。
设计要点:
rivetkit的Ctx<A>方法一律作为ActorContext的薄透传(thin pass-through),不承载核心业务逻辑;新加的封装方法只在边界处做 CBOR 编解码并委托给 core。因此使用rivetkit编写的代码,其底层行为由rivetkit-core保证。
二、设计约束总览
CLAUDE.md 用五条硬性约束划定了 SDK 的功能边界。下面逐条展开,并给出源码证据与实战含义。
1. 不提供vars临时变量 API
TypeScript SDK 中存在ctx.vars(临时变量 API),其存在原因是 TS 用户状态存放在框架中。而 Rust 中:
- 临时状态(ephemeral state):直接作为 Actor struct 的普通字段,例如 chat-room-rust 中
ChatRoom的started_at_ms: i64; - 持久化状态(persisted state):即关联类型
Actor::State,由框架负责保存。
因此不要在rivetkit中添加vars访问器来镜像 TypeScript。Rust 开发者应依赖 Rust 的结构体所有权模型,将生命周期内数据放在 struct 字段上,将需要跨休眠/恢复的数据放在State中。
2. 主 API 为Actor+ 每个 action 一个Handles<M>实现
封装层的主入口 API 是traitActor+ 针对每个 action 的Handles<M>实现,生命周期与分发逻辑集中在run_actor中:
run_actor位于 src/start.rs,它一次性完成启动阶段(输入解码、状态创建、create、on_create、on_start)并将失败原样上报给运行时握手,随后进入事件循环分发。- 用户 Actor 应通过
register_actor/register_actor_with注册(见 src/registry.rs),而不是直接编写事件循环。
trait 化模型下,Action 的完整定义如下(src/action.rs):
pub trait Action: serde::Serialize + DeserializeOwned + Send + Sync + 'static { type Output: serde::Serialize + DeserializeOwned + Send + 'static; const NAME: &'static str; }Handles<A>trait(src/action.rs)要求为每个 Action 实现一个返回Future的handle方法:
pub trait Handles<A: Action>: Actor + Sized { type Future: Future<Output = Result<A::Output>> + Send + 'static; fn handle(self: Arc<Self>, ctx: Ctx<Self>, action: A) -> Self::Future; }Actor::Actions关联类型使用元组(tuple)声明可处理的 action 集合,宏为从 1 到 128 个元数(TUPLE_ARITY_MAX = 128,见 src/action.rs)自动生成ActionSet实现。分发时按name匹配,并用encode_positional/decode_positional在边界完成 CBOR 的位置参数编解码。
3. 旧事件循环 API 保留兼容,新代码一律使用 trait API
Registry::register、Registry::register_with、Start、RuntimeEvent构成的旧式事件循环 API仍然可用,但已在源码中标记为#[deprecated]:
- src/registry.rs 对
register的弃用说明为 "use register_actor/register_actor_with and implement Actor + Handles instead"。 - src/start.rs 中的
Start<A>携带ctx、input、is_new、snapshot、hibernated、events等字段,Events<A>::recv()返回RuntimeEvent<A>,供手动事件循环消费。
新示例与文档应使用 trait API,事件循环仅作为迁移期的兼容路径。对应地,trybuildUI 测试(如 action_set_missing_handle.rs 及其.stderr快照)会编译失败并给出明确提示,确保漏掉Handles实现的编译期错误可读、可诊断。
4. 状态修改自动标记 dirty,request_save仅用于显式保存点
这是最容易影响持久化正确性的规则。源码 src/context.rs 的实现细节如下:
Ctx::state_mut()(src/context.rs):先设置dirty标志为true,再返回写锁保护的可变引用;StateMut的Drop(src/context.rs):在释放写锁后自动调用inner.request_save(...),无需用户手动操作;Ctx::set_state()(src/context.rs):整体替换状态并立即触发request_save。
因此:
- 通过
state_mut()/set_state()修改的持久化状态会自动标记为 dirty 并触发保存; request_save()(src/context.rs)只应在需要显式保存点(例如希望状态尽快落盘、或需要配合request_save_with_opts)时调用;- 事件循环中收到
SerializeState事件时(src/start.rs),运行时先判断state_dirty(),若 dirty 则先执行on_state_changehook,再编码StateDelta(CBOR)并清除 dirty 标志。
实战建议:不要在每次字段修改后额外调用
request_save,这会造成冗余的保存请求;依赖 dirty 标记机制,只在需要显式持久化时机时使用。
5. 不支持 workflow 事件:工作流引擎归属rivetkit-typescript
Rust Actor永远不会托管 workflow,工作流引擎由rivetkit-typescript拥有。对应地:
- 不要在 Rust SDK 中暴露
ActorEvent::WorkflowHistoryRequested/WorkflowReplayRequested的Event变体或类型; - 在
Event::from_core(src/event.rs)中对这类 core 变体使用unreachable!处理; - src/start.rs 中,若 core 侧仍然投递了这两个事件,则返回
ActorRuntime::NotConfigured错误("workflow history" / "workflow replay")。
该约束保证 Rust SDK 的职责边界清晰:Rust 处理有状态 Actor、action、队列、定时与连接,workflow 编排交给 TypeScript 生态。
三、Actortrait 全貌:生命周期、HTTP、WebSocket 与连接
Actortrait 定义在 src/actor.rs,它通过 12 个关联类型和约 15 个可覆写方法定义了 Actor 的完整行为面。关联类型如下:
| 关联类型 | 含义 | 说明 |
|---|---|---|
State | 持久化状态 | 需Serialize + DeserializeOwned,由框架保存 |
Input | 创建输入 | 需DeserializeOwned + Default,缺失时用默认值 |
Actions | action 集合 | 实现ActionSet<Self>的元组 |
Events | 可广播事件集合 | 实现EventSet的元组 |
Queue | 队列消息集合 | 实现QueueSet<Self>的元组 |
ConnParams | 连接参数 | 需DeserializeOwned + Default |
ConnState | 连接状态 | 需Serialize + DeserializeOwned + Clone |
Action | 分发时的 Action 类型 | 通常为action::Raw |
关键常量与方法:
const HAS_DATABASE: bool:声明 Actor 是否使用用户数据库。注册时 src/registry.rs 会将其并入ActorConfig::has_database;同时,内部存储(state、KV、queue)始终使用 SQLite,因此在未编译sqlite-localfeature 时,actor_config会强制remote_sqlite = true,即通过 engine 路由 SQLite。const CONCURRENT_HTTP_CALLBACKS: bool = false+MAX_CONCURRENT_HTTP_CALLBACKS: usize = 128+MAX_CONCURRENT_LIVE_HTTP_CALLBACK_STARTS: usize = 0:开启并发 HTTP 回调时,run_actor会为 Standard 与 Live 两类回调建立独立的信号量池(src/start.rs),池满时回复 HTTP 429。classify_http_request(&Request) -> HttpCallbackClass:可将请求分类到 Standard / Live 回调池。admit_http_request(...) -> HttpCallbackAdmission:允许 Actor 在回调构造/注册前同步执行准入逻辑,返回Untracked/Tracked { release }/Reject。Tracked的release被包装在HttpCallbackAdmissionGuard中,在回调回复交给 core 或回调被 drop 时自动执行(src/start.rs)。- 生命周期 hook:
create_state、create、on_create、on_start、run、on_state_change、on_sleep、on_destroy。 - 连接 hook:
create_conn_state、on_before_connect、on_connect、on_disconnect、on_subscribe。 - HTTP / WebSocket:
on_fetch(返回Response)、on_fetch_response(返回ActorHttpResponse,默认为 buffered,可返回StreamingResponse实现流式响应)、on_websocket(默认bail!("websockets not supported"))。
所有 hook 默认实现均为空操作或保守默认值,最小化 Actor 只需实现create_state/create(否则返回ActorRuntime::NotConfigured)。
运行时启动流程
Registry::start()(src/registry.rs)是独立二进制入口,它会阻塞直到收到 SIGINT/SIGTERM,然后取消并排空。运行模式由环境变量RIVETKIT_RUNTIME_MODE决定:
Envoy(默认):持有一个长生命周期出站 envoy,start_envoy直至信号后排空;Serverless:启动 HTTP listener,在首个请求时惰性启动并缓存 envoy,对应into_serverless_runtime。
连接设置(包括用于派生或复用本地 engine 的RIVET_ENGINE_BINARY_PATH)从环境变量读取。由于 Rust 没有隐式运行时保持进程存活,Registry::start会阻塞;需要编程式生命周期控制(测试、嵌入场景)时,应使用serve/serve_with_config并自行驱动CancellationToken。
四、状态、输入与快照:Start结构解析
run_actor消费的Start<A>(src/start.rs)由 core 的ActorStart通过wrap_start包装而来,包含:
ctx: Ctx<A>:类型化上下文;input: Input<A>:启动输入。is_present()判断是否携带字节;decode()/decode_or_default()/decode_or(f)解码 CBOR 输入;缺失输入时回退到A::Input::default()(与 TypeScriptcreateState(undefined)语义一致)。若输入缺失且类型不是(),则返回ActorRuntime::MissingInput。is_new: bool:是否为新建 Actor;snapshot: Snapshot:休眠恢复的快照。is_new()判断是否为全新实例;decode::<S>()解码状态快照,空快照返回None;hibernated: Vec<Hibernated<A>>:休眠时保留的连接(携带连接状态);events: Events<A>:事件流,recv()/try_recv()返回RuntimeEvent<A>。
启动阶段(src/start.rs)作为一个整体可失败单元执行:若有快照则解码恢复状态,否则调用create_state;随后A::create、is_new时的on_create、以及on_start。整个阶段通过startup_readyoneshot channel 与运行时握手,失败会以真实原因(而非通道关闭的泛化错误)上报。
RuntimeEvent<A>(src/event.rs)共有 10 个变体:Action、Http、QueueSend、WebSocketOpen、ConnOpen、ConnClosed、Subscribe、SerializeState、Sleep、Destroy。值得注意的是Events::recv会内部消化ConnectionOpen、DisconnectConn、RunWake这类运行时握手事件(src/start.rs),只向用户暴露业务相关事件。
五、Ctx<A>能力地图:状态、SQL、调度、广播与客户端
Ctx<A>(src/context.rs)内部持有ActorContext+StateCell(值 + dirty 标志)+ 惰性Client+ 可选的当前ConnCtx。核心能力分类如下:
状态访问:state()(只读)、state_mut()(写 + 自动 dirty + Drop 时 request_save)、set_state()、state_dirty()、clear_state_dirty()、set_initial_state()。
SQLite 访问:sql()返回SqliteDb;db_exec/db_query/db_execute/db_run提供 CBOR 边界的便捷方法。SqliteDbExt::transaction(src/sqlite.rs)提供 commit-on-success 事务助手,支持命名事务与超时。
调度与定时:schedule()支持after(延迟)、at(定点)、cancel、get、list;cron()支持set(cron 表达式 + timezone + max_history)、every(间隔 Duration)、get/list/delete/history。
唤醒控制:keep_awake(future)(对应 TSctx.waitUntil,future 在途期间 Actor 不会休眠)、keep_awake_region()(返回 guard,drop 时释放)、abort_signal()/aborted()(运行时销毁信号)、register_task(注册运行时拥有的后台任务,shutdown 时与优雅期限竞争)。
事件与广播:broadcast(name, event)/emit(E)(向连接广播事件)、conns()/conns_vec()/disconnect_conn/disconnect_conns。
生命周期控制:sleep()(请求休眠)、destroy()、stop_with_error(message)(带错误停止,engine 记录为 crash 并应用崩溃处理)、set_alarm。
类型化客户端:client()惰性构造rivetkit_client::Client(Bare 编码 + WebSocket 传输),用于跨 Actor 调用;TypedClientExt(src/typed_client.rs)提供get_typed/get_or_create_typed返回TypedActorHandle<A>,可在编译期绑定 Actor 类型后安全调用其 action。
连接上下文ConnCtx<A>:id()、params()、state()/set_state()、send(name, event)、disconnect(reason)、is_hibernatable()。连接参数与状态同样以 CBOR 编解码。
注:
kv()已被标记为 deprecated,建议使用嵌入式 SQLite(sql())或 Actor state 替代(src/context.rs);set_prevent_sleep/prevent_sleep同样已弃用为 no-op,改用keep_awake或wait_until。
六、队列:Queue与QueueSet
队列 API 与 action 保持同样的 trait 化结构(src/queue.rs):
QueueMessage:Serialize + DeserializeOwned,含NAME与Reply类型;HandlesQueue<M>:每个消息一个handle_queue实现;QueueSet<A>:元组组合(支持最多 16 个消息),按名称分发;Ctx::queue()返回Queue<'_>,发送助手会 CBOR 编码消息体,*_raw变体原样透传字节。
事件循环投递RuntimeEvent::QueueSend,其中wait/timeout_ms字段支持阻塞等待语义;src/start.rs 中队列处理器不存在时返回ActorRuntime::NotFound。
七、状态持久化的代码级验证
约束 4(自动 dirty)与约束 2(run_actor生命周期)可由源码直接验证:
state_mut()在返回写守卫前dirty.store(true, Ordering::Release);StateMut::drop先drop(guard)释放锁,再request_save—— 注释明确指出这镜像了 TypeScript 的 write-through 状态代理,不依赖优雅关机过程来保证持久化;SerializeState事件处理:dirty 时先跑on_state_change,再encode_state_delta(CBOR 编码StateDelta::ActorState),最后clear_state_dirty。
src/persist.rs 进一步提供state_delta/state_deltas/conn_hibernation_delta/conn_hibernation_removed_delta等显式 delta 构造工具,供手动保存场景使用。
八、测试与示例:从 e2e 到 UI 编译测试
rivetkit的测试体系覆盖多个层级:
- In-process e2e 测试:tests/test_harness_e2e.rs 通过
test::setup(registry)(src/test.rs)派生或复用本地 engine(解析顺序:显式路径 →RIVET_ENGINE_BINARY_PATH→ 健康 engine → 本地工作区构建 → 缓存二进制 → 可选验证下载),为并发测试分配唯一 pool 名避免跨路由,随后用类型化 handle 发送 action 并断言往返结果。 - 模块单元测试:src/start.rs 内置
LifecycleActor测试覆盖完整生命周期(create_state→create→on_create→on_start→run→on_sleep/on_destroy),以及连接预检、订阅拒绝、并发 HTTP 回调的 429 与 admission 释放计数。 - trybuild UI 测试:tests/trybuild.rs 与 tests/ui/ 中的
.rs+.stderr快照验证漏写Handles/HandlesQueue实现等场景的编译错误信息质量。 - 官方示例:examples/chat-room-rust 是完整的 trait API 演示:定义
ChatRoomActor(含HAS_DATABASE = true、on_fetch默认 404 之外的 SQLite 建表、sendMessage/getHistory/getStats三个 action、newMessage事件广播、ctx.state_mut()计数),其 main.rs 仅两行:example_chat_room_rust::registry().start().await。
九、快速上手与最佳实践总结
最小可运行结构(参考 test_harness_e2e.rs):
use rivetkit::{Action, Actor, Ctx, Handles, Registry, action}; use serde::{Deserialize, Serialize}; #[derive(Serialize, Deserialize)] struct Echo { value: String } impl Action for Echo { type Output = String; const NAME: &'static str = "echo"; } struct MyActor; impl Actor for MyActor { type State = (); type Input = (); type Actions = (Echo,); type Events = (); type Queue = (); type ConnParams = (); type ConnState = (); type Action = action::Raw; // 需要 SQLite 用户数据库时:const HAS_DATABASE: bool = true; } impl Handles<Echo> for MyActor { type Future = std::pin::Pin<Box<dyn std::future::Future<Output = anyhow::Result<String>> + Send>>; fn handle(self: std::sync::Arc<Self>, _ctx: Ctx<Self>, action: Echo) -> Self::Future { Box::pin(async move { Ok(action.value) }) } } #[tokio::main] async fn main() -> anyhow::Result<()> { let mut registry = Registry::new(); registry.register_actor::<MyActor>("myActor"); registry.start().await // 或 test::setup(registry) 进行进程内测试 }关键规则速查:
- 新代码一律使用
Actor+Handlestrait API,旧事件循环 API 仅为兼容保留; - 临时数据放 struct 字段,持久化数据放
State,不要期望varsAPI; - 依赖
state_mut()/set_state()的自动 dirty 保存,request_save()仅在需要显式保存点时使用; Ctx方法是薄透传,不要在应用层绕过类型化 API 直接操作 core;- Rust Actor 不托管 workflow,相关工作流能力在
rivetkit-typescript侧; - 测试优先使用
test::setup进程内 harness;需要控制生命周期时用serve+CancellationToken而非阻塞的start(); - 未编译
sqlite-localfeature 时内部存储默认走远程 SQLite(经 engine 路由),用户数据库需显式声明HAS_DATABASE。
遵循以上设计约束,即可写出与rivetkit-core运行时语义一致、可休眠、可持久化、可测试的 Rust Actor 应用。
【免费下载链接】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),仅供参考