diff --git a/server-rs/crates/api-server/src/app.rs b/server-rs/crates/api-server/src/app.rs
index 390031ef9..732e0d046 100644
--- a/server-rs/crates/api-server/src/app.rs
+++ b/server-rs/crates/api-server/src/app.rs
@@ -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
{
+ 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,
@@ -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"]
diff --git a/server-rs/crates/platform-llm/src/observability_tests.rs b/server-rs/crates/platform-llm/src/observability_tests.rs
index 16cb87dfc..a166c9b7b 100644
--- a/server-rs/crates/platform-llm/src/observability_tests.rs
+++ b/server-rs/crates/platform-llm/src/observability_tests.rs
@@ -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(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");
}