如何在Rust自定义Tokio-Tracing订阅器中避免Hyper无限日志循环
说明
我已查阅《在tokio tracing订阅器中抑制外部库日志》和《如何关闭其他 crate 发出的追踪事件?》,但它们无法解决我的问题。
背景
我正在开发一个Tokio tracing订阅器,用于捕获由tracing instrumentation的库生成的日志,并通过hyper crate导出到后端服务。但hyper同样基于Tokio tracing进行了 instrumentation,导致其生成的日志被我的订阅器捕获并导出,进而引发无限日志循环。
请注意,集成该订阅器的应用也可能使用hyper crate。我仅希望抑制订阅器实现内部使用hyper所产生的追踪事件,而非应用层使用hyper产生的事件。如链接问题中建议的添加hyper过滤器会抑制所有hyper事件,无论其使用场景。以下是简化的可运行代码及问题细节:
/* hyper = { version = "0.14.7", features = [ "full" ] } tokio = { version = "1.33.0", features = ["full"] } tracing = "0.1.25" tracing-core = "0.1" pin-project = "1.1.3" tracing-subscriber = "0.3.17" */ use pin_project::pin_project; use std::future::Future; use std::pin::Pin; use std::task::{Context, Poll}; use tracing::info; use tracing::Event; use tracing::Metadata; use tracing::Subscriber; use tracing_core::span::Id; use tracing_core::span::Record; // Define a task-local variable for suppression tokio::task_local! { static SUPPRESSED: bool; } struct SimpleSubscriber; impl Subscriber for SimpleSubscriber { fn enabled(&self, _metadata: &Metadata<'_>) -> bool { !is_logging_suppressed() } fn new_span(&self, _: &tracing::span::Attributes<'_>) -> Id { Id::from_u64(10) } fn event(&self, event: &Event<'_>) { if is_logging_suppressed() { return; } // Extract the required metadata before moving into the async context let target = event.metadata().target().to_string(); let name = event.metadata().name().to_string(); let level = event.metadata().level().clone(); let suppressed_future = SUPPRESSED.scope(true, async move { // the actual logging event, the hyper tracing event generated here should be suppressed hyper_wrapper::make_hyper_call("mock_upload_event", &target, &name, &level).await; }); let suppress_logging_future = SuppressLogging { inner: suppressed_future, }; // Use tokio::spawn instead of spawn_local tokio::task::spawn(suppress_logging_future); } fn record(&self, _: &Id, _: &Record<'_>) {} fn record_follows_from(&self, _: &Id, _: &Id) {} fn enter(&self, _: &Id) {} fn exit(&self, _: &Id) {} } #[pin_project] struct SuppressLogging<F> { #[pin] inner: F, } impl<F: Future> Future for SuppressLogging<F> { type Output = F::Output; fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> { let this = self.project(); this.inner.poll(cx) } } fn is_logging_suppressed() -> bool { SUPPRESSED.try_with(|val| *val).unwrap_or(false) } #[tokio::main] async fn main() { let subscriber = SimpleSubscriber; tracing::subscriber::set_global_default(subscriber) .expect("Setting default subscriber failed."); // Make a small Hyper call in the main method lib::some_func().await; } // Define the module within the same file mod lib { use tracing::info; use tracing_core::Level; pub async fn some_func() { // this should be logged info!(name: "this-is-also-logged", "test"); // the tracing event generated by hyper crate should not be suppressed. super::hyper_wrapper::make_hyper_call("main method", "main", "startup", &Level::INFO).await; } } mod hyper_wrapper { use hyper::Client; use hyper::Uri; use tracing_core::Level; pub async fn make_hyper_call(context: &str, target: &str, name: &str, level: &Level) { let client = Client::new(); let uri = "http://www.google.com".parse::<Uri>().unwrap(); let response = client.get(uri).await; match response { Ok(_response) => { //let _body = hyper::body::to_bytes(response).await.unwrap(); println!( "{} - [{}] - {} - {} ", context, level, target, name, ); } Err(e) => { println!("{} - [{}] - {} - {} - Error: {:?}", context, level, target, name, e); } } } }
问题说明
- 我的Tokio tracing订阅器通过hyper将日志导出到后端,但hyper自身的追踪事件会被订阅器捕获,引发无限循环;
- 直接过滤所有hyper日志会影响应用层或其他依赖库的hyper日志;
- 我尝试使用tokio::task_local设置抑制标记,但该标记无法在hyper内部的异步调用链中传播,导致抑制失效。
核心问题
如何仅过滤订阅器实现内部hyper(及reqwest、tonic等同类crate)产生的日志,同时不影响应用层或其他地方的hyper日志?
示例场景说明
SimpleSubscriber是我将发布为crate的订阅器;main函数属于集成该crate的应用,同时使用lib模块;lib模块的日志及其中使用hyper产生的日志应被订阅器导出;- 订阅器内部通过
event()调用make_hyper_call时产生的hyper日志需被抑制,避免循环; - 尝试的tokio::task_local方案因标记无法传播而失效。
解决方案
要解决这个问题,核心是区分订阅器内部发起的异步任务和应用层的任务,确保抑制逻辑能准确覆盖内部操作,同时不干扰外部任务。以下是几种可行方案:
1. 基于tracing span的上下文判断
在订阅器内部发起导出操作时,创建一个专属标记span,然后在enabled方法中检查当前是否处于该span内,同时判断事件是否来自需要抑制的crate:
impl Subscriber for SimpleSubscriber { fn enabled(&self, metadata: &Metadata<'_>) -> bool { // 检查当前是否处于订阅器导出的span中 let is_in_export_span = tracing::Span::current().name() == "subscriber_export"; // 定义需要抑制的目标crate let is_suppressed_crate = metadata.target().starts_with("hyper") || metadata.target().starts_with("reqwest") || metadata.target().starts_with("tonic"); // 仅当处于导出span且是目标crate时才抑制 !(is_in_export_span && is_suppressed_crate) } fn event(&self, event: &Event<'_>) { if !self.enabled(event.metadata()) { return; } let target = event.metadata().target().to_string(); let name = event.metadata().name().to_string(); let level = event.metadata().level().clone(); tokio::spawn(async move { // 创建专属span标记导出操作,span会自动在异步链中传播 let export_span = tracing::info_span!("subscriber_export"); let _guard = export_span.enter(); hyper_wrapper::make_hyper_call("mock_upload_event", &target, &name, &level).await; }); } // 其他方法保持不变... }
这种方式无需依赖任务本地变量,利用tracing自身的span上下文传播特性,可靠性更高。
2. 修复任务本地变量的传播问题
如果坚持使用tokio::task_local,需要确保所有内部异步操作都在标记的scope内执行,同时去掉多余的SuppressLogging包装器(它不会帮助传播上下文):
fn event(&self, event: &Event<'_>) { if is_logging_suppressed() { return; } let target = event.metadata().target().to_string(); let name = event.metadata().name().to_string(); let level = event.metadata().level().clone(); // 直接在spawn的任务中设置scope,确保所有内部异步调用都继承该标记 tokio::spawn(async move { SUPPRESSED.scope(true, async move { hyper_wrapper::make_hyper_call("mock_upload_event", &target, &name, &level).await; }).await; }); }
如果需要使用tokio::spawn_local(保留线程本地上下文),需确保任务是!Send类型,或者通过tokio::sync::oneshot传递结果。
3. 结合任务ID标记
可以在订阅器内部生成任务时记录任务ID,然后在enabled方法中检查当前任务ID是否属于订阅器内部任务:
use std::sync::atomic::{AtomicU64, Ordering}; static INTERNAL_TASK_ID: AtomicU64 = AtomicU64::new(0); impl Subscriber for SimpleSubscriber { fn enabled(&self, metadata: &Metadata<'_>) -> bool { let is_suppressed_crate = metadata.target().starts_with("hyper") || metadata.target().starts_with("reqwest") || metadata.target().starts_with("tonic"); if !is_suppressed_crate { return true; } // 检查当前任务是否是订阅器内部任务 let current_task_id = tokio::task::Id::current().as_u64(); let internal_task_id = INTERNAL_TASK_ID.load(Ordering::Relaxed); current_task_id != internal_task_id } fn event(&self, event: &Event<'_>) { if !self.enabled(event.metadata()) { return; } let target = event.metadata().target().to_string(); let name = event.metadata().name().to_string(); let level = event.metadata().level().clone(); let task = tokio::spawn(async move { hyper_wrapper::make_hyper_call("mock_upload_event", &target, &name, &level).await; }); // 记录内部任务ID INTERNAL_TASK_ID.store(task.id().as_u64(), Ordering::Relaxed); } // 其他方法保持不变... }
这种方式适合简单场景,但需要注意任务ID的复用问题,可结合任务完成后的清理逻辑优化。
内容的提问来源于stack exchange,提问作者LKB

