ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

深入解析 Rerun 的 re_test_mocks:进程内 OTLP 与 PostHog 测试桩的设计与实战

深入解析 Rerun 的 re_test_mocks:进程内 OTLP 与 PostHog 测试桩的设计与实战 深入解析 Rerun 的 re_test_mocks进程内 OTLP 与 PostHog 测试桩的设计与实战【免费下载链接】rerunVisualize, query, and stream to train on multimodal robotics data.项目地址: https://gitcode.com/GitHub_Trending/re/rerunRerun 仓库当前工作区根目录GitHub_Trending/re/rerun的 crates/tests/re_test_mocks/README.md 描述了一组专为测试而生的进程内服务替身in-process server doublesMockOtlpCollector与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 可以看到支撑这一切的依赖组合tonicgRPC 服务端启用router、transport、gzip、opentelemetry-proto启用了gen-tonic与trace特性以生成 proto 类型、axumHTTP 服务端、tokio/tokio-stream异步运行时与 listener 流、parking_lot无异步上下文的互斥锁以及serde_jsonPostHog JSON 解析。二、crate 结构无根导出显式子模块README 明确说明crate 根不 re-export 任何东西消费者必须直接深入子模块引用use re_test_mocks::otlp::MockOtlpCollector; use re_test_mocks::posthog::MockPostHog;这从 src/lib.rs 得到印证——它只声明了三个pub modassert、otlp、posthog。显式子模块路径让每个 mock 的依赖关系清晰独立otlp依赖 tonic/protoposthog依赖 axum/serde_json也避免了把两种协议的类型混在同一个命名空间里。assert子模块导出的是伴生宏assert_sink_empty!见下文第五节与两个 mock 配合完成无多余流量的收尾断言。三、MockOtlpCollector进程内 OTLP TraceServiceMockOtlpCollector的核心职责是接收 OTLPExportRPC把嵌套的ResourceSpans → ScopeSpans → spans结构展平flatten成一个个独立 span连同请求的 gRPC metadata 一起存入共享缓冲区。3.1 展平后的观测粒度ReceivedSpan一次携带 N 个 span 的ExportRPC 会产生 N 个ReceivedSpan每个都复制同一份metadata但拥有各自的resource/scope/spanpub struct ReceivedSpan { pub metadata: tonic::metadata::MetadataMap, pub resource: OptionResource, pub scope: OptionInstrumentationScope, pub span: Span, }这是测试断言的基本单元wait_for的谓词一次只看一个 span匹配时只弹出那一个同批次的兄弟 span 留在缓冲区等待后续匹配。源码注释特别强调了一个 OTLP 与 PostHog 的行为差异ReceivedSpan没有BadRequest变体——因为 tonic 在进入 handler 之前就完成了 proto 反序列化任何畸形Export都会在到达 mock 之前以InvalidArgument被拒绝src/otlp.rs。相比之下 PostHog 是明文 JSON必须自己处理坏请求见第四节。3.2 生命周期与关键 APIMockOtlpCollector的完整公开面如下src/otlp.rsAPI签名说明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()- VecReceivedSpan缓冲区快照克隆取出缓冲区不动is_empty()- bool缓冲区是否为空初始、完全消费或clear之后clear()- ()清空缓冲区wait_for(predicate, timeout)async fn - ResultReceivedSpan, OtlpWaitTimeout等待并消费下一个满足谓词的 span见 3.3shutdown()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传播到每个展平 spanrecords_request_metadata_on_every_flattened_span自定义 gRPC metadata如x-test-tag复制到每个展平 spanwait_for_pops_only_the_matched_span谓词命中b时a、c留在缓冲区wait_for_returns_when_matching_span_arrives50ms 延迟到达的 span 被wait_for捕获并验证总耗时远小于 5 秒预算通知驱动、非轮询wait_for_timeout_leaves_buffer_untouched超时后快照里能看到keptspan缓冲区仍然持有它assert_sink_empty_panics_with_diagnostic缓冲区非空时宏以expected empty, got 1 request(s)恐慌shutdown_completes_cleanlyshutdown().await后任务干净结束。四、MockPostHog进程内 PostHog /batch 端点MockPostHog与MockOtlpCollector刻意保持形状一致——同样的spawn/addr/endpoint/received/is_empty/clear/wait_for/shutdown表面src/posthog.rs让两种协议的测试读起来完全一致。区别在协议层。4.1 展平 /batch 数组ReceivedEventaxumhandler 只路由根路径/的 POST其他方法/路径得到 axum 默认的 405/404 且不记录见测试rejects_non_post_methods。每个 POST 的/batch数组被展平为每个条目一个ReceivedEventpub struct ReceivedEvent { pub headers: HeaderMap, pub event: EventBody, }其中EventBody有两种形态pub enum EventBody { Ok(serde_json::Value), // 解析成功的 batch 条目 BadRequest { raw: Vecu8, error: String }, // 坏请求诊断 }无法解析的 JSON或解析成功但没有/batch数组的请求会作为一个BadRequest事件入缓冲区handler 返回400——协议违规被显式暴露而不是静默接受EventBody::as_ok()返回OptionValue适合在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) - VecT其中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 无 TLSendpoint()返回明文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),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进