news 2026/9/13 3:21:24

SpacetimeDB 过程并发行为测试模块解析:sdk-test-procedure-concurrency 的设计与验证

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
SpacetimeDB 过程并发行为测试模块解析:sdk-test-procedure-concurrency 的设计与验证

SpacetimeDB 过程并发行为测试模块解析:sdk-test-procedure-concurrency 的设计与验证

【免费下载链接】SpacetimeDBDevelopment at the speed of light项目地址: https://gitcode.com/GitHub_Trending/sp/SpacetimeDB

SpacetimeDB 的 Procedure(过程)与 Reducer(归约器)是两类不同的服务端执行单元,当过程在执行中途挂起(sleep_until)、普通归约器或调度归约器同时被触发时,二者的执行顺序与穿插行为直接影响应用的数据一致性。本文围绕仓库中的 modules/sdk-test-procedure-concurrency/README.md 及其源码 src/lib.rs,深入解析这个专门用于隔离验证过程并发行为的 Rust 测试模块:它为什么独立存在、如何通过表行插入顺序观测并发穿插、底层sleep_until的宿主 ABI 实现,以及 SDK 测试套件如何从客户端侧验证这些行为。读完本文,你将理解 SpacetimeDB 中 Procedure 与 Reducer 并发调度的真实语义,并掌握这套"以数据行顺序观测并发"的测试方法论。

模块定位:为什么要有一个专门的并发测试模块

原文档(modules/sdk-test-procedure-concurrency/README.md)明确了该模块的职责:

This module isolates procedure concurrency behavior that currently only has Rust module coverage.

也就是说,这个模块是独立隔离过程并发行为的测试载体,而这类行为当前只有 Rust 模块实现覆盖。它与其姊妹模块sdk-test-procedure分离的原因在原文中有直接说明:

It is separate fromsdk-test-procedureso the shared procedure test suite can continue targeting other module languages without also requiringctx.sleep_untilsupport.

翻译过来即:共享的 Procedure 测试套件要继续面向其他模块语言(TypeScript、C#、C++)做交叉验证,而这些语言并不都支持ctx.sleep_until。因此需要把依赖ctx.sleep_until的并发场景单独拆出来,放进一个"仅 Rust 模块"的测试模块,避免拖累多语言共享套件。

从工作区配置也能印证这一分工:根目录 Cargo.toml 将本模块纳入 workspace 成员,模块自身的 Cargo.toml 声明为:

[package] name = "sdk-test-procedure-concurrency-module" version = "0.1.0" edition.workspace = true license-file = "LICENSE" [lib] crate-type = ["cdylib"] [dependencies] log.workspace = true [dependencies.spacetimedb] workspace = true features = ["unstable"]

关键点有两个:一是crate-type = ["cdylib"],表明它编译为 Wasm 动态库供 SpacetimeDB 宿主加载;二是启用spacetimedbcrate 的unstablefeature,因为ProcedureContextsleep_untilScheduleAt等过程能力当前仍属不稳定 API。

并发行为的观测载体:ProcedureConcurrencyRow 表

过程与归约器的并发穿插,在这个模块里不是靠时间戳或日志断言,而是靠表行的插入顺序来客观记录。核心表定义在 src/lib.rs:

#[table(public, accessor = procedure_concurrency_row)] struct ProcedureConcurrencyRow { #[auto_inc] insertion_order: u32, insertion_context: String, }
  • insertion_order#[auto_inc]自增列,由数据库自动分配严格递增的序号。同一表内行号的大小关系,就等价于这些写入在不同执行单元间发生的先后次序,这是整个测试设计的基石。
  • insertion_context:写入来源的标记字符串,用于区分该行是谁插入的("procedure_before""reducer""procedure_after""scheduled_reducer""scheduled_procedure_before""scheduled_procedure_after")。

统一入口函数 insert_procedure_concurrency_row 在给定事务上下文内完成插入:

fn insert_procedure_concurrency_row(ctx: &TxContext, insertion_context: &str) { ctx.db.procedure_concurrency_row().insert(ProcedureConcurrencyRow { insertion_order: 0, insertion_context: insertion_context.into(), }); }

insertion_order传 0 即可,#[auto_inc]会由引擎覆盖为真实递增序号。

核心机制:ProcedureContext、with_tx 与 sleep_until

要理解这些测试场景,必须先弄清楚 Procedure 与 Reducer 的执行模型差异。在 crates/bindings/src/lib.rs 中,ProcedureContext被定义为一个#[non_exhaustive]结构体,保存调用方Identity、过程启动时间timestamp等信息,且过程必须以&mut ProcedureContext作为第一个参数

该上下文暴露两个关键方法(见 crates/bindings/src/lib.rs 与 crates/bindings/src/lib.rs):

pub fn sleep_until(&mut self, timestamp: Timestamp) { let new_time = sys::procedure::sleep_until(timestamp.to_micros_since_unix_epoch()); let new_time = Timestamp::from_micros_since_unix_epoch(new_time); self.timestamp = new_time; } pub fn with_tx<T>(&mut self, body: impl Fn(&TxContext) -> T) -> T { with_tx(body, self.sender(), self.connection_id()) }
  • with_tx:在过程内部开启一个读写事务执行body。注意过程的数据库操作必须显式包在with_tx中,且文档提醒body可能被多次执行(重试语义),闭包内不应写入外部可变状态。
  • sleep_until:把过程挂起直到指定时刻,返回时self.timestamp已更新为唤醒后的新时间。这正是过程能够"在两次插入之间让出执行权、等待其他执行单元介入"的能力来源——归约器(Reducer)没有这个能力,它一次事务执行完即结束。

从宿主端看,sleep_until的底层是一条异步 ABI 导入。在 crates/core/src/host/wasm_common.rs 中,procedure_sleep_untilprocedure_http_request一起被列为异步链接($link_async!,模块spacetime_10.3)。更值得注意的实现细节在 crates/core/src/host/wasmtime/wasmtime_module.rs:对于同步Wasmtime 实例,procedure_sleep_until会被替换为一个直接报错的桩函数:

fn procedure_sleep_until_sync_stub(_: Caller<'_, WasmInstanceEnv>, _: i64) -> anyhow::Result<i64> { anyhow::bail!("procedure_sleep_until is only available in async instances") }

这从源码层面解释了原文档的措辞:ctx.sleep_until依赖异步实例支持,并非所有模块语言/运行配置都具备,这正是它必须从多语言共享测试套件中剥离出来的根本原因。

轮询辅助工具:poll_until_tx_true

由于过程挂起后需要"等某个外部事件发生再继续",模块实现了一个基于sleep_until的轮询辅助函数(src/lib.rs):

#[derive(Copy, Clone, Debug)] struct PollOptions { timeout: Duration, poll_interval: Duration, } impl Default for PollOptions { fn default() -> Self { Self { timeout: Duration::from_secs(10), poll_interval: Duration::from_millis(100), } } } fn poll_until_tx_true(ctx: &mut ProcedureContext, pred: impl Fn(&TxContext) -> bool, options: PollOptions) { let deadline = ctx.timestamp + options.timeout; log::info!("poll_until_tx_true: will give up at {deadline}"); while ctx.timestamp < deadline { let try_again = ctx.timestamp + options.poll_interval; log::info!("poll_until_tx_true: sleeping until {try_again}"); ctx.sleep_until(try_again); if ctx.with_tx(&pred) { log::info!("poll_until_tx_true: succeeded, returning now"); return; } log::info!("poll_until_tx_true: false"); } panic!("poll_until_tx_true: exceeded timeout {:?}", options.timeout) }

行为逻辑:每poll_interval(默认 100ms)用sleep_until睡到下一个检查点,然后开一个事务执行谓词pred;一旦谓词为真立即返回,超过timeout(默认 10s)则panic!。整个过程通过log::info!输出详细进度,便于在SPACETIME_LOG中排查。

测试场景一:Procedure 与 Reducer 的穿插(interleaving)

第一个核心场景是 procedure_sleep_between_inserts:

#[procedure] fn procedure_sleep_between_inserts(ctx: &mut ProcedureContext) { ctx.with_tx(|ctx| insert_procedure_concurrency_row(ctx, "procedure_before")); poll_until_tx_true( ctx, |tx| { tx.db .procedure_concurrency_row() .iter() .any(|row| row.insertion_context != "procedure_before") }, Default::default(), ); ctx.with_tx(|ctx| insert_procedure_concurrency_row(ctx, "procedure_after")); }

执行流程:

  1. 过程先插入一行insertion_context = "procedure_before"
  2. 进入轮询,等待任何非"procedure_before"的行出现——即等待其他执行单元写入;
  3. 一旦出现,过程插入"procedure_after"

配合普通的归约器 insert_reducer_row(只插入一行"reducer"),整个测试的断言目标就是:归约器必须在过程的两次插入之间执行,最终表内顺序为procedure_before < reducer < procedure_after。若归约器被阻塞到过程结束后才运行,顺序会变成procedure_before < procedure_after < reducer,测试即失败。

这个场景验证的正是"过程挂起期间,其他客户端触发的归约器可以穿插执行"这一并发语义。

测试场景二:Procedure 与调度归约器(Scheduled Reducer)的并发

第二个场景引入了调度机制。先看调度归约器表(src/lib.rs):

#[table(accessor = scheduled_reducer_row, scheduled(insert_scheduled_reducer))] struct ScheduledReducerRow { #[primary_key] #[auto_inc] scheduled_id: u64, scheduled_at: ScheduleAt, } #[reducer] fn insert_scheduled_reducer(ctx: &ReducerContext, _schedule: ScheduledReducerRow) { ctx.db.procedure_concurrency_row().insert(ProcedureConcurrencyRow { insertion_order: 0, insertion_context: "scheduled_reducer".into(), }); }

#[table(scheduled(...))]声明了一个"每插入一行即调度一次"的归约器:插入ScheduledReducerRow就相当于在scheduled_at时刻排队执行insert_scheduled_reducer。而过程 procedure_schedule_reducer_between_inserts 把"插入 before 行"与"调度归约器"放在同一个事务里,然后同样轮询等待非"procedure_before"行出现,最后插入"procedure_after"

#[procedure] fn procedure_schedule_reducer_between_inserts(ctx: &mut ProcedureContext) { ctx.with_tx(|ctx| { insert_procedure_concurrency_row(ctx, "procedure_before"); ctx.db.scheduled_reducer_row().insert(ScheduledReducerRow { scheduled_id: 0, scheduled_at: ctx.timestamp.into(), }); }); // ... 与场景一相同的轮询 ... ctx.with_tx(|ctx| insert_procedure_concurrency_row(ctx, "procedure_after")); }

注意scheduled_at: ctx.timestamp.into()——调度的触发时刻就是过程启动时刻。期望的行为是:过程挂起后,调度归约器在procedure_beforeprocedure_after之间执行,最终顺序为procedure_before < scheduled_reducer < procedure_after

测试场景三:调度过程与调度归约器"不穿插"(已知行为)

第三个场景展示了当前引擎的已知限制,源码注释毫不避讳地写明了这一点(src/lib.rs):

#[table(accessor = scheduled_procedure_row, scheduled(scheduled_procedure_sleep_between_inserts))] struct ScheduledProcedureRow { #[primary_key] #[auto_inc] scheduled_id: u64, scheduled_at: ScheduleAt, } #[procedure] fn scheduled_procedure_sleep_between_inserts(ctx: &mut ProcedureContext, _schedule: ScheduledProcedureRow) { ctx.with_tx(|ctx| insert_procedure_concurrency_row(ctx, "scheduled_procedure_before")); // Unfortunately, we can't poll and wake on event here, // as (until the related upstream issue is fixed) // the scheduled reducer actually won't run until after this procedure fully completes. ctx.sleep_until(ctx.timestamp + Duration::from_secs(10)); ctx.with_tx(|ctx| insert_procedure_concurrency_row(ctx, "scheduled_procedure_after")); }

调度过程插入"scheduled_procedure_before"后直接睡 10 秒,再插入"scheduled_procedure_after"。注释明确解释:在相关上游问题修复之前,调度归约器实际要等该过程完全结束后才会运行,因此这里无法像前两个场景那样用"事件唤醒"式轮询,只能靠固定时长睡眠。

启动这一切的入口归约器 schedule_procedure_then_reducer 同时调度两件事:

#[reducer] fn schedule_procedure_then_reducer(ctx: &ReducerContext) { ctx.db.scheduled_procedure_row().insert(ScheduledProcedureRow { scheduled_id: 0, scheduled_at: ctx.timestamp.into(), }); ctx.db.scheduled_reducer_row().insert(ScheduledReducerRow { scheduled_id: 0, scheduled_at: (ctx.timestamp + Duration::from_secs(2)).into(), }); }
  • 调度过程在ctx.timestamp(立即)执行;
  • 调度归约器在ctx.timestamp + 2s执行。

由于过程的睡眠窗口是 10 秒,归约器(2 秒后)会在过程仍在挂起时就已经到期。但按当前引擎行为,它仍须等过程结束。因此该场景断言的行序是scheduled_procedure_before < scheduled_procedure_after < scheduled_reducer——一个"不穿插"的负向验证,用于把现有(可能并不理想的)调度器语义固化下来,防止无意间改变。

SDK 测试套件:从客户端侧验证并发行为

模块本身只提供服务端行为,真正的验收在 Rust SDK 测试套件中完成。在 sdks/rust/tests/test.rs 中,rust_procedure_concurrency测试模块将上述场景一一对应为四个测试用例:

测试函数客户端子命令对应模块场景
procedure_reducer_interleavingprocedure-reducer-interleaving场景一:过程与归约器穿插
procedure_reducer_same_client_not_interleavedprocedure-reducer-same-client-interleaved场景一 + 同客户端限制
procedure_concurrent_with_scheduled_reducerprocedure-concurrent-with-scheduled-reducer场景二:过程与调度归约器并发
scheduled_procedure_scheduled_reducer_not_interleavedscheduled-procedure-scheduled-reducer-not-interleaved场景三:调度器单执行槽(不穿插)

每个用例都通过platform_test_builder绑定本模块并生成私有项绑定(.with_generate_private_items(true))。第四个用例的注释(sdks/rust/tests/test.rs)进一步道明了设计意图:

Test that the scheduler has only a single active execution slot, which can be occupied by a long-running or suspended procedure. We're not attached to this behavior, and in fact it should be changed. At that time, this test should be altered to demonstrate that the execution is interleaved.

即:调度器当前只有一个活动执行槽,可被长时间运行或被挂起的过程占用;团队并不认可该行为,未来修复后此测试应改为验证穿插执行。

对应的测试客户端位于 sdks/rust/tests/procedure-concurrency-client,其 README.md 说明该客户端的目标是"测试当前 Rust 模块 ABI 与 Rust SDK 的归约器/过程并发行为"。客户端主程序 src/main.rs 从命令行参数取测试名、从SPACETIME_SDK_TEST_DB_NAME取数据库名,再由 src/test_handlers.rs 分发到各执行函数。

客户端侧的验证手法很有参考价值(见 test_handlers.rs):

  • 订阅SELECT * FROM procedure_concurrency_row;全表;
  • on_insert回调里按insertion_context分派到状态机字段,记录各自拿到的insertion_order
  • 三者齐备后断言严格顺序。例如场景一在 test_handlers.rs 中断言before < reducer < after,场景二断言before < scheduled_reducer < after,场景三则断言before < after < scheduled_reducer(test_handlers.rs)。

此外测试还通过TestCounter协调多个异步回调的完成(add_test/wait_for_all),并用Arc<Mutex<...>>保护共享观测状态,保证每个回调只被报告一次(ordering_checked标志 +take())。

如何运行

本模块的行为验证统一走 Rust SDK 测试套件。按 sdk-test-procedure 的 README 的做法,在仓库根目录执行:

cargo test -p spacetimedb-sdk procedure

即可运行名称含procedure的用例,其中就包括上述rust_procedure_concurrency模块中的四个并发测试。测试框架会负责启动独立 SpacetimeDB 实例、编译并发布本模块(sdk-test-procedure-concurrency-module)、生成绑定代码并运行客户端。

总结

sdk-test-procedure-concurrency是一个小而精的专项测试模块,它回答了一个关键问题:当过程在sleep_until挂起时,普通归约器与调度归约器到底能不能穿插执行?模块用ProcedureConcurrencyRow表的自增insertion_order把并发时序"物化"为可断言的顺序数据,配合poll_until_tx_true的轮询机制,在模块层完成了三种场景的观测——归约器可穿插、调度归约器可穿插、调度过程与调度归约器暂不可穿插(已知行为,待上游修复)。这套"用数据行顺序固化并发语义"的测试方法论,以及"把依赖不稳定 API 的用例从多语言共享套件中剥离"的工程取舍,对任何需要精确控制服务端执行顺序的 SpacetimeDB 模块开发都有直接的借鉴意义。

【免费下载链接】SpacetimeDBDevelopment at the speed of light项目地址: https://gitcode.com/GitHub_Trending/sp/SpacetimeDB

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/13 3:17:39

SVM支持向量机原理与实战:从鸢尾花理解决策边界与核函数

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/13 3:12:46

嵌入式最小硬件系统全解析:从电源时钟到PCB调试

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/13 3:12:16

3 分钟上手 PDF 书签修复、合并与提图:PDF 补丁丁工具箱实战指南

3 分钟上手 PDF 书签修复、合并与提图&#xff1a;PDF 补丁丁工具箱实战指南 【免费下载链接】PDFPatcher PDF补丁丁——PDF工具箱&#xff0c;可以编辑书签、剪裁旋转页面、解除限制、提取或合并文档&#xff0c;探查文档结构&#xff0c;提取图片、转成图片等等 项目地址: …

作者头像 李华