修掉并发测试下 HTTP/Provider span 偶发采集为空的 tracing interest 竞态

- 根因:span!/info_span! 宏在 callsite 的缓存 interest 为 never 时会静默返回空 span
  (tracing-0.1.44/src/macros.rs 的 span! 分支),而 DefaultCallsite::register 只在调用点
  首次被命中时计算一次 interest,且当进程里只注册过一个 dispatcher 时会退化成
  dispatcher::get_default()——也就是命中线程自己的 dispatcher(tracing-core-0.1.36
  callsite.rs 的 Rebuilder::JustOne)。libtest 默认并发下,没有 subscriber 的普通测试线程
  一旦抢到 http.request / llm.request 调用点的首次注册,就会把它永久缓存成 never,
  于是 CI 偶发看到 0 个 span(app.rs:603、observability_tests.rs:98)
- 关键:set_default 与 with_subscriber 在 interest 缓存这件事上**等价**(都只是新建 Dispatch
  并触发一次 rebuild_interest,只能纠正“已经注册过”的调用点,纠正不了 JustOne→get_default
  这条首次注册分支),所以真正起作用的是**在自家 subscriber 下命中同一个调用点做热身**,
  把首次注册的顺序握在自己手里;register_callsite 覆写只用于避免本 subscriber 触发的重建
  把其它调用点永久标记成 never
- api-server app::tests::http_tracing:不再用 with_subscriber,改为在请求前用同一
  http.request 调用点热身到连续两轮采集成功,并在整个请求期间持有 scoped default
  (已注明 set_default 是线程绑定,只适用于默认的 current_thread #[tokio::test])
- platform-llm observability_tests:新增 run_under_capture 固定整段流程的 scoped default,
  并新增 warm_up_provider_span_callsite(用指向刚释放 loopback 端口的最小失败请求命中同一
  llm.request 调用点,连续两轮采集成功才继续)
- 断言未放宽:仍要求每个被拒绝请求恰好 1 个 HTTP span、每次 Provider 调用恰好 1 个
  llm.request span;未使用 sleep
- 本地实跑:cargo test -p api-server --bin api-server app::tests::http_tracing(默认并发与
  --test-threads=1 各 20 次全绿)、app::tests:: 91 用例并发 10 次全绿、--skip bgfilter_worker
  --skip wallet_refund_outbox 的 1133 用例全量 bin 2 次全绿;cargo test -p platform-llm --lib
  (162 passed,含 3 条观测用例)连跑 20 次全绿
This commit is contained in:
2026-10-04 00:56:55 +08:00
parent 05707e61e4
commit 977b2f0c75
2 changed files with 134 additions and 18 deletions
+61 -6
View File
@@ -510,11 +510,12 @@ mod tests {
use tracing::{
Metadata, Subscriber,
field::{Field, Visit},
instrument::WithSubscriber,
span::{Attributes, Id, Record},
subscriber::Interest,
};
use super::*;
use crate::app::make_http_request_span;
use crate::state::{BackpressureState, HttpRequestPermitPoolKind};
#[derive(Clone, Default)]
@@ -534,6 +535,18 @@ mod tests {
}
impl Subscriber for HttpSpanCapture {
/// tracing 的 callsite interest 是**进程级**缓存,只在调用点首次命中时计算一次。
/// 这里只为 HTTP span 声明 always:任何以本 subscriber 参与的重建都会让
/// `http.request` 调用点保持可用;其余调用点返回 sometimes,交给 `enabled`
/// 在真正创建 span 时判断,避免把其它调用点永久标记成 never。
fn register_callsite(&self, metadata: &Metadata<'_>) -> Interest {
if metadata.is_span() && metadata.name() == "http.request" {
Interest::always()
} else {
Interest::sometimes()
}
}
fn enabled(&self, metadata: &Metadata<'_>) -> bool {
metadata.is_span() && metadata.name() == "http.request"
}
@@ -557,6 +570,41 @@ mod tests {
fn exit(&self, _: &Id) {}
}
fn http_span_warmup_request() -> Request<Body> {
Request::builder()
.uri("/__genarrative_http_span_interest_warmup__")
.body(Body::empty())
.expect("warmup request should build")
}
/// 让 `http.request` 调用点在本进程内稳定可用,而不是依赖调用点首次注册时命中线程的 dispatcher。
///
/// `DefaultCallsite::register` 只在调用点首次被命中时计算一次 interest,计算时会用
/// `DISPATCHERS.rebuilder()`;当进程里只有一个 dispatcher 时它会退化成
/// `dispatcher::get_default()`,也就是**命中线程自己的** dispatcher。并发跑测试时,
/// 没有 subscriber 的线程一旦抢到这次注册,`http.request` 调用点就会被永久缓存成
/// `never`,之后 `span!` 宏会静默跳过 span 创建,测试便会观察到 0 个 span。
///
/// 这里不 sleep 赌时序:每轮都用 `set_default` 在我们的 subscriber 下安装 scoped default
/// (内部 `Dispatch::new` 会触发 tracing 重建 interest 缓存),并直接命中同一个
/// `http.request` 调用点做探测;连续两轮都探测到 span 才认为缓存已经稳定。
fn ensure_http_request_span_interest() {
let mut captured_rounds = 0;
for _ in 0..8 {
let probe = HttpSpanCapture::default();
{
let _guard = tracing::subscriber::set_default(probe.clone());
let _ = make_http_request_span(&http_span_warmup_request());
}
let captured = !probe.0.lock().expect("span capture should lock").is_empty();
captured_rounds = if captured { captured_rounds + 1 } else { 0 };
if captured_rounds >= 2 {
return;
}
}
panic!("无法在测试 subscriber 下恢复 http.request 调用点的 interest 缓存");
}
async fn assert_rejection_observed(
app: Router,
request: Request<Body>,
@@ -567,11 +615,18 @@ mod tests {
let expected_request_id = request.headers().get("x-request-id").cloned();
let path = request.uri().path().to_string();
let capture = HttpSpanCapture::default();
let response = app
.oneshot(request)
.with_subscriber(capture.clone())
.await
.expect("rejected request should complete");
// 先把 http.request 调用点的进程级 interest 缓存校正到当前 subscriber(见上方注释),
// 再在整个请求期间持有 scoped default。注意 `set_default` 是**线程绑定**的:这里成立
// 是因为本模块用的是默认 `#[tokio::test]`(current-thread 运行时,请求与断言线程一致);
// 一旦改成 multi_thread / spawn 到别的线程,必须换成 `with_current_subscriber()` 或
// `with_subscriber`,否则请求会在没有 dispatcher 的线程上跑。
ensure_http_request_span_interest();
let response = {
let _guard = tracing::subscriber::set_default(capture.clone());
app.oneshot(request)
.await
.expect("rejected request should complete")
};
assert_eq!(response.status(), expected_status);
let request_id = response.headers()["x-request-id"]
@@ -11,7 +11,6 @@ use tokio::sync::oneshot;
use tracing::{
Instrument, Subscriber,
field::{Field, Visit},
instrument::WithSubscriber,
span::{Attributes, Id, Record},
};
use tracing_subscriber::{Layer, layer::Context, prelude::*, registry::LookupSpan};
@@ -222,8 +221,70 @@ fn request() -> LlmRunRequest {
.with_model("requested-model")
}
/// 在整段流程期间把 capture 固定为当前 dispatcher。
///
/// tracing 的 callsite interest 是**进程级**缓存,且只在调用点首次被命中时计算一次
/// (`DefaultCallsite::register` → `DISPATCHERS.rebuilder()`;进程里只有一个 dispatcher 时会
/// 退化成命中线程自己的 dispatcher)。并发跑测试时,本文件之外那些没有 subscriber 的测试线程
/// 只要抢到 `llm.request` 调用点的首次注册,它就会被永久缓存成 `never`,`span!` 宏随后静默
/// 跳过 span 创建,`provider_span` 便会看到 0 个 span。
///
/// 注意:`set_default` 与 `with_subscriber` 在 interest 缓存这件事上**等价**——两者都只是新建一个
/// `Dispatch` 并触发一次 `rebuild_interest`,只能纠正“已经注册过”的调用点,纠正不了
/// `Rebuilder::JustOne → get_default() → 缓存 never` 这条首次注册分支;真正把调用点救回来的是
/// `warm_up_provider_span_callsite` 在自家 subscriber 下命中同一调用点(覆盖两种顺序,
/// 并且把 dispatcher 从“只在每次 poll 生效”变成整段流程生效)。
async fn run_under_capture<F: Future>(capture: &Capture, future: F) -> F::Output {
// `set_default` 是线程绑定的:本文件用的是默认 `#[tokio::test]`(current-thread 运行时,
// 被测 future 与断言在同一条线程上)。若以后改成 multi_thread 或把流程 spawn 出去,
// 必须换回 `with_current_subscriber()` / `with_subscriber`,否则 span 会在没有 dispatcher 的线程上丢。
let _guard =
tracing::subscriber::set_default(tracing_subscriber::registry().with(capture.clone()));
future.await
}
/// 让 `llm.request` 调用点在本进程内稳定可用:用一个必然失败的最小请求(指向刚释放的
/// loopback 端口,连接直接被拒绝)命中同一个调用点,直到连续两轮都能采集到 span 为止。
/// 请求失败无妨——只要 `llm.request` span 被创建就说明调用点在我们的 subscriber 下注册成功。
/// 不 sleep、不放宽断言。
async fn warm_up_provider_span_callsite() {
let listener = TcpListener::bind("127.0.0.1:0").expect("warmup listener should bind");
let address = listener
.local_addr()
.expect("warmup address should resolve");
drop(listener);
let config = LlmConfig::new(
LlmProvider::OpenAiCompatible,
format!("http://{address}"),
PRIVATE_KEY.into(),
"default-model".into(),
500,
0,
1,
)
.expect("warmup config should build");
let client = LlmClient::new(config).expect("warmup client should build");
let mut captured_rounds = 0;
for _ in 0..8 {
let probe = Capture::default();
{
let _guard = tracing::subscriber::set_default(
tracing_subscriber::registry().with(probe.clone()),
);
let _ = client.run(request()).await;
}
let captured = !probe.spans.lock().expect("capture should lock").is_empty();
captured_rounds = if captured { captured_rounds + 1 } else { 0 };
if captured_rounds >= 2 {
return;
}
}
panic!("无法在测试 subscriber 下恢复 llm.request 调用点的 interest 缓存");
}
#[tokio::test]
async fn provider_span_covers_awaited_execution_and_keeps_parent_without_arguments() {
warm_up_provider_span_callsite().await;
let fixture = provider_fixture(
"200 OK",
"application/json",
@@ -231,7 +292,7 @@ async fn provider_span_covers_awaited_execution_and_keeps_parent_without_argumen
);
let capture = Capture::default();
let response =
async {
run_under_capture(&capture, async {
async {
let mut future = Box::pin(fixture.client.run(request()));
assert!(capture.spans.lock().unwrap().values().all(|span| span.name != "llm.request"));
@@ -245,8 +306,7 @@ async fn provider_span_covers_awaited_execution_and_keeps_parent_without_argumen
fixture.release.send(()).unwrap();
future.await.unwrap()
}.instrument(tracing::info_span!("test.request")).await
}
.with_subscriber(tracing_subscriber::registry().with(capture.clone()))
})
.await;
fixture.server.join().unwrap();
assert_eq!(response.text, "completed");
@@ -255,6 +315,7 @@ async fn provider_span_covers_awaited_execution_and_keeps_parent_without_argumen
#[tokio::test]
async fn stream_callbacks_inherit_provider_span_and_keep_result() {
warm_up_provider_span_callsite().await;
let fixture = provider_fixture(
"200 OK",
"text/event-stream",
@@ -263,7 +324,7 @@ async fn stream_callbacks_inherit_provider_span_and_keep_result() {
let capture = Capture::default();
let mut deltas = Vec::new();
let response =
async {
run_under_capture(&capture, async {
async {
let mut future = Box::pin(fixture.client.stream_run(request(), |delta| {
tracing::info!(target: "llm_observability_test_delta", "delta received");
@@ -276,8 +337,7 @@ async fn stream_callbacks_inherit_provider_span_and_keep_result() {
fixture.release.send(()).unwrap();
future.await.unwrap()
}.instrument(tracing::info_span!("test.request")).await
}
.with_subscriber(tracing_subscriber::registry().with(capture.clone()))
})
.await;
fixture.server.join().unwrap();
assert_eq!(response.text, "hello");
@@ -291,21 +351,21 @@ async fn stream_callbacks_inherit_provider_span_and_keep_result() {
#[tokio::test]
async fn provider_span_preserves_upstream_errors_and_closes_on_cancellation() {
warm_up_provider_span_callsite().await;
let fixture = provider_fixture(
"401 Unauthorized",
"application/json",
r#"{"error":{"message":"upstream-rejected"}}"#,
);
let capture = Capture::default();
let error = async {
let error = run_under_capture(&capture, async {
async {
fixture.release.send(()).unwrap();
fixture.client.run(request()).await.unwrap_err()
}
.instrument(tracing::info_span!("test.request"))
.await
}
.with_subscriber(tracing_subscriber::registry().with(capture.clone()))
})
.await;
fixture.server.join().unwrap();
assert!(
@@ -315,7 +375,7 @@ async fn provider_span_preserves_upstream_errors_and_closes_on_cancellation() {
let fixture = provider_fixture("200 OK", "application/json", "{}");
let capture = Capture::default();
async {
run_under_capture(&capture, async {
async {
let mut future = Box::pin(fixture.client.run(request()));
tokio::select! {
@@ -326,7 +386,8 @@ async fn provider_span_preserves_upstream_errors_and_closes_on_cancellation() {
drop(future);
fixture.release.send(()).unwrap();
}.instrument(tracing::info_span!("test.request")).await
}.with_subscriber(tracing_subscriber::registry().with(capture.clone())).await;
})
.await;
fixture.server.join().unwrap();
capture.assert_completed("run");
}