
SpacetimeDB 过程并发行为测试模块解析sdk-test-procedure-concurrency 的设计与验证【免费下载链接】SpacetimeDBDevelopment at the speed of light项目地址: https://gitcode.com/GitHub_Trending/sp/SpacetimeDBSpacetimeDB 的 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因为ProcedureContext、sleep_until、ScheduleAt等过程能力当前仍属不稳定 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.rspub 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_txT(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_until与procedure_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::Resulti64 { 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)); }执行流程过程先插入一行insertion_context procedure_before进入轮询等待任何非procedure_before的行出现——即等待其他执行单元写入一旦出现过程插入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_before与procedure_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 cant poll and wake on event here, // as (until the related upstream issue is fixed) // the scheduled reducer actually wont 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. Were 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_reducertest_handlers.rs。此外测试还通过TestCounter协调多个异步回调的完成add_test/wait_for_all并用ArcMutex...保护共享观测状态保证每个回调只被报告一次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),仅供参考