深入解析 Rerun 的 re_test_mocks:进程内 OTLP 与 PostHog 测试桩的设计与实战
【免费下载链接】rerunVisualize, query, and stream to train on multimodal robotics data.项目地址: https://gitcode.com/GitHub_Trending/re/rerun
Rerun 仓库(当前工作区根目录GitHub_Trending/re/rerun)的 crates/tests/re_test_mocks/README.md 描述了一组专为测试而生的进程内服务替身(in-process server doubles):MockOtlpCollector与MockPostHog。它们分别以真实 wire protocol 完整实现 OTLP gRPCTraceService与 PostHog/batchHTTP 接口,运行在测试进程内的临时端口上,捕获每一条出站遥测请求,并通过**通知驱动(notification-driven)**的wait_for(…)/received()访问器取代轮询式断言。本文将结合 crate 源码、内置测试与真实调用方re_perf_telemetry的用法,完整还原这两个 mock 的设计理念、核心 API、并发语义与实战写法,帮助你为任何需要"捕获出站遥测流量"的 Rust 测试写出同样的基础设施。
一、为什么需要进程内服务替身
在测试 Rerun 这类会向外部服务发送遥测(OTel trace)与产品分析(PostHog)数据的代码时,一个经典难题是:不能真的把测试流量发到生产后端,而 mock 掉整个客户端又会让被测试的序列化、gRPC 拦截器、批处理等真实链路失去验证价值。
re_test_mocks的答案介于两者之间:启动一个真实的服务器进程,跑真实的协议实现,但不离开测试进程。
- 对 OTLP,它用
tonic启动一个真正的 gRPCTraceService服务端(src/otlp.rs); - 对 PostHog,它用
axum启动一个真正的 HTTP 服务端,处理POST /上的/batch批量捕获端点(src/posthog.rs)。
两者都绑定到 OS 分配的127.0.0.1临时端口,因此并行测试之间天然隔离、零冲突。被测试代码只需把 exporter / client 指向endpoint()返回的地址,就能走完全部真实协议路径——只是终点是内存中的断言缓冲区。
从 Cargo.toml 可以看到支撑这一切的依赖组合:tonic(gRPC 服务端,启用router、transport、gzip)、opentelemetry-proto(启用了gen-tonic与trace特性以生成 proto 类型)、axum(HTTP 服务端)、tokio/tokio-stream(异步运行时与 listener 流)、parking_lot(无异步上下文的互斥锁)以及serde_json(PostHog JSON 解析)。
二、crate 结构:无根导出,显式子模块
README 明确说明:crate 根不 re-export 任何东西,消费者必须直接深入子模块引用:
use re_test_mocks::otlp::MockOtlpCollector; use re_test_mocks::posthog::MockPostHog;这从 src/lib.rs 得到印证——它只声明了三个pub mod:assert、otlp、posthog。显式子模块路径让每个 mock 的依赖关系清晰独立(otlp依赖 tonic/proto,posthog依赖 axum/serde_json),也避免了把两种协议的类型混在同一个命名空间里。
assert子模块导出的是伴生宏assert_sink_empty!(见下文第五节),与两个 mock 配合完成"无多余流量"的收尾断言。
三、MockOtlpCollector:进程内 OTLP TraceService
MockOtlpCollector的核心职责是:接收 OTLPExportRPC,把嵌套的ResourceSpans → ScopeSpans → spans结构展平(flatten)成一个个独立 span,连同请求的 gRPC metadata 一起存入共享缓冲区。
3.1 展平后的观测粒度:ReceivedSpan
一次携带 N 个 span 的ExportRPC 会产生 N 个ReceivedSpan,每个都复制同一份metadata,但拥有各自的resource/scope/span:
pub struct ReceivedSpan { pub metadata: tonic::metadata::MetadataMap, pub resource: Option<Resource>, pub scope: Option<InstrumentationScope>, pub span: Span, }这是测试断言的基本单元:wait_for的谓词一次只看一个 span,匹配时只弹出那一个,同批次的兄弟 span 留在缓冲区等待后续匹配。
源码注释特别强调了一个 OTLP 与 PostHog 的行为差异:ReceivedSpan没有BadRequest变体——因为 tonic 在进入 handler 之前就完成了 proto 反序列化,任何畸形Export都会在到达 mock 之前以InvalidArgument被拒绝(src/otlp.rs)。相比之下 PostHog 是明文 JSON,必须自己处理坏请求(见第四节)。
3.2 生命周期与关键 API
MockOtlpCollector的完整公开面如下(src/otlp.rs):
| API | 签名 | 说明 |
|---|---|---|
spawn() | async fn spawn() -> Self | 绑定127.0.0.1:0并开始服务;只有服务器任务真正开始执行后才返回,避免后续请求与任务首次 poll 竞争 |
addr() | -> SocketAddr | 实际监听地址 |
endpoint() | -> String | 形如http://127.0.0.1:PORT,可直接用于 tonic / OTLP exporter 配置;仅支持明文 gRPC,无 TLS——若调用方错误地套上 HTTPS 客户端会连接失败,这是测试配置 bug,而非 mock 的缺陷 |
received() | -> Vec<ReceivedSpan> | 缓冲区快照(克隆取出,缓冲区不动) |
is_empty() | -> bool | 缓冲区是否为空(初始、完全消费或clear之后) |
clear() | -> () | 清空缓冲区 |
wait_for(predicate, timeout) | async fn -> Result<ReceivedSpan, OtlpWaitTimeout> | 等待并消费下一个满足谓词的 span(见 3.3) |
shutdown() | async fn (self) | 优雅关闭:发 shutdown 信号并等待服务器任务结束 |
服务端还支持gzip 压缩:TraceServiceServer::new(service).accept_compressed(Gzip).send_compressed(Gzip)(src/otlp.rs),与真实 collector 的压缩行为对齐。
Drop是 fire-and-forget 的:直接 drop 会发 shutdown 信号但不会等待任务结束(tokio 会 detach 任务让它跑完)。需要在运行时拆除前确保 in-flight 响应处理完毕的测试,必须显式调用shutdown().await。
3.3 wait_for:通知驱动的消费式等待
wait_for是整个 crate 的设计精髓——测试不需要 sleep 轮询。其实现(src/otlp.rs)是经典的"先武装通知、再扫描缓冲区"循环:
let deadline = tokio::time::Instant::now() + timeout; loop { // 1. 先武装 waiter,再扫描——避免 push 恰好发生在 scan 与 await 之间而被错过 let notified = self.state.notify.notified(); tokio::pin!(notified); notified.as_mut().enable(); { let mut buffer = self.state.received.lock(); if let Some(pos) = buffer.iter().position(&predicate) { return Ok(buffer.remove(pos)); // 2. 命中即弹出(消费) } } // 3. 未命中则等待通知或超时 if tokio::time::timeout_at(deadline, notified).await.is_err() { return Err(OtlpWaitTimeout { snapshot: self.received() }); } }值得注意的语义细节:
- 命中即移除:谓词匹配的 span 按到达顺序被
remove弹出;同批次兄弟留在原地。因此连续多次wait_for可以逐个消费匹配项。 - 超时不破坏缓冲区:超时时返回
OtlpWaitTimeout,其中snapshot是超时瞬间缓冲区内容的克隆,缓冲区本身原封不动——后续wait_for可以继续基于同一缓冲区等待。 - 消费点即断言点:正如源码注释所说,"wait_for 弹出之后缓冲区里剩下的东西才是真正的多余流量"——一条意外重传、一个测试忘记断言的 span,都会被
assert_sink_empty!逮住。 - 处理器先记录后响应:
exporthandler 在返回响应之前就已把展平结果写入缓冲区并notify_waiters()(src/otlp.rs)。因此客户端c.export(…).await返回 Ok 时,这批 span 必然已在缓冲区中——简单的同步received()断言即可,不必wait_for。
3.4 源码自带的测试佐证
src/otlp.rs 内置了一组高质量测试,可直接当作 API 使用范例:
records_one_received_span_per_proto_span:一次携带 3 个 span 的 Export 展平为 3 个缓冲项,且resource传播到每个展平 span;records_request_metadata_on_every_flattened_span:自定义 gRPC metadata(如x-test-tag)复制到每个展平 span;wait_for_pops_only_the_matched_span:谓词命中"b"时,"a"、"c"留在缓冲区;wait_for_returns_when_matching_span_arrives:50ms 延迟到达的 span 被wait_for捕获,并验证总耗时远小于 5 秒预算(通知驱动、非轮询);wait_for_timeout_leaves_buffer_untouched:超时后快照里能看到"kept"span,缓冲区仍然持有它;assert_sink_empty_panics_with_diagnostic:缓冲区非空时宏以expected empty, got 1 request(s)恐慌;shutdown_completes_cleanly:shutdown().await后任务干净结束。
四、MockPostHog:进程内 PostHog /batch 端点
MockPostHog与MockOtlpCollector刻意保持形状一致——同样的spawn/addr/endpoint/received/is_empty/clear/wait_for/shutdown表面(src/posthog.rs),让两种协议的测试读起来完全一致。区别在协议层。
4.1 展平 /batch 数组:ReceivedEvent
axumhandler 只路由根路径/的 POST(其他方法/路径得到 axum 默认的 405/404 且不记录,见测试rejects_non_post_methods)。每个 POST 的/batch数组被展平为每个条目一个ReceivedEvent:
pub struct ReceivedEvent { pub headers: HeaderMap, pub event: EventBody, }其中EventBody有两种形态:
pub enum EventBody { Ok(serde_json::Value), // 解析成功的 batch 条目 BadRequest { raw: Vec<u8>, error: String }, // 坏请求诊断 }- 无法解析的 JSON,或解析成功但没有
/batch数组的请求,会作为一个BadRequest事件入缓冲区,handler 返回400——协议违规被显式暴露,而不是静默接受; EventBody::as_ok()返回Option<&Value>,适合在wait_for谓词中使用;expect_parsed()在遇到BadRequest时恐慌并打印原始 body 与错误,适合直接断言场景(src/posthog.rs)。
4.2 与 OTLP mock 的语义对齐
PostHog 的wait_for、超时类型PosthogWaitTimeout、"先记录后响应"(handler 在返回响应前写入缓冲区,见 src/posthog.rs)等设计与 OTLP 完全一致。服务端通过axum::serve(listener, app).with_graceful_shutdown(...)优雅关闭。
内置测试(src/posthog.rs)覆盖了:单 batch 多事件展平、请求头复制到每个展平事件、畸形 JSON → 400 + BadRequest、缺/batch→ 400 + 错误信息含/batch、非 POST 方法 → 405 且不记录、谓词弹出、延迟到达、超时、clear/is_empty状态、expect_parsed恐慌与shutdown。
五、assert_sink_empty!:无流量收尾断言
assert_sink_empty!宏(src/assert.rs)是wait_for的"镜像":前者断言"等到过",后者断言"不该有的都没有"。
#[macro_export] macro_rules! assert_sink_empty { ($sink:expr $(,)?) => {{ let __received = $sink.received(); assert!( __received.is_empty(), "expected empty, got {} request(s):\n{:#?}", __received.len(), __received, ); }}; }- 它对任何暴露
fn received(&self) -> Vec<T>(其中T: Debug)的类型都有效,因此天然同时适用于MockOtlpCollector与MockPostHog,甚至自定义 sink; - 必须传引用:
assert_sink_empty!(&collector); - 失败时打印缓冲区完整内容(
{:#?}),直接展示"意外到达了什么",与wait_for的消费语义形成闭环——wait_for弹掉预期流量后,宏确认没有多余流量。
六、真实调用方:re_perf_telemetry 的端到端认证测试
re_test_mocks并非孤立设施。在 crates/utils/re_perf_telemetry/src/telemetry.rs 的authed_exporter_sends_bearer_metadata测试中,它被用来验证一条完整的JWT 认证导出管线:
use re_test_mocks::otlp::MockOtlpCollector; let collector = MockOtlpCollector::spawn().await; let jwt = re_auth::Jwt::try_from(TEST_JWT.to_owned()).unwrap(); let provider = super::Arc::new(StaticCredentialsProvider::new(jwt)); // 把 mock 的 endpoint 喂给真实 exporter 构建器 let exporter = super::build_rerun_authed_span_exporter_with_provider(&collector.endpoint(), provider) .unwrap(); let tracer_provider = SdkTracerProvider::builder() .with_span_processor(BatchSpanProcessor::builder(exporter).build()) .build(); let tracer = tracer_provider.tracer("test"); // 发一个 span,强制 flush(不等默认 5 秒调度延迟) { let span = tracer.start("authed_test_span"); drop(span); } tracer_provider.force_flush().ok(); // 通知驱动等待,随即校验 gRPC metadata 里的 authorization 头 let received = collector .wait_for(|_| true, Duration::from_secs(10)) .await .expect("collector should receive at least one span"); let auth = received .metadata .get("authorization") .expect("authorization metadata missing") .to_str() .expect("authorization should be ASCII"); assert_eq!(auth, format!("Bearer {TEST_JWT}"));这个用例展示了 mock 的核心价值场景:被测试代码走的是真实的 exporter → tonic interceptor → gRPC metadata 完整链路,而 mock 从 wire 层把authorization: Bearer <jwt>元数据捕获出来供断言——这是任何"假客户端"都做不到的端到端验证。测试里还体现了两个实用技巧:用force_flush()触发批处理导出以跳过默认调度延迟,以及用wait_for(|_| true, 10s)表达"等待任意 span 到达"。
七、实践要点与限制总结
7.1 推荐测试模式
综合 README、源码注释与真实调用方,一个典型的re_test_mocks测试可以组织为:
let collector = MockOtlpCollector::spawn().await; // 或 MockPostHog::spawn() // 配置被测代码指向 collector.endpoint()(http://127.0.0.1:PORT) // 触发一次导出…… // 1. 同步断言:handler 先记录后响应,客户端 await 完成后数据必在缓冲区 assert_eq!(collector.received().len(), 1); // 2. 通知驱动消费:按谓词弹出预期流量(超时不破坏缓冲区,可重试) let span = collector .wait_for(|s| s.span.name == "expected", Duration::from_secs(5)) .await .unwrap(); // 3. 无流量收尾:弹完后不应有任何多余流量 assert_sink_empty!(&collector); // 4. fire-and-forget 请求场景:显式优雅关闭,确保 in-flight 响应处理完毕 collector.shutdown().await;7.2 已知限制(以源码为准)
- OTLP mock 无 TLS:
endpoint()返回明文http://地址,仅支持明文 gRPC。套 HTTPS 客户端会连接失败(属于测试配置问题,见 src/otlp.rs 注释)。 - PostHog mock 只路由
/的 POST:其他方法或路径返回 405/404 且不记录。 Drop是 fire-and-forget:需要优雅关闭验证时必须显式shutdown().await。- 缓冲区无容量上限:不消费的话请求会持续累积——这正是
assert_sink_empty!存在的意义。
八、总结
re_test_mocks用约 600 行代码为 Rerun 的测试体系提供了一个可复用的"进程内遥测黑盒"范式:用真实协议栈启动临时服务器、把 wire 数据展平成可断言的观测粒度、用通知驱动取代轮询、用"等待 + 清空断言"的双宏闭环覆盖"该有的有、不该有的没有"两种验证需求。无论你是想理解 Rerun 的遥测测试如何工作,还是计划为自己的项目搭建 OTLP / PostHog 捕获型测试基础设施,crates/tests/re_test_mocks/ 都是值得直接复用的参考实现。
【免费下载链接】rerunVisualize, query, and stream to train on multimodal robotics data.项目地址: https://gitcode.com/GitHub_Trending/re/rerun
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考