use super::*;
use std::sync::{Arc, Mutex, OnceLock};
use tracing::Level;
use tracing_subscriber::layer::{Context, Layer, SubscriberExt};
use tracing_subscriber::Registry;
type Sink = Arc<Mutex<Vec<String>>>;
static ACTIVE_SINK: Mutex<Option<Sink>> = Mutex::new(None);
#[derive(Clone, Default)]
struct WarnCapture;
impl<S: tracing::Subscriber> Layer<S> for WarnCapture {
fn on_event(&self, event: &tracing::Event<'_>, _ctx: Context<'_, S>) {
struct Visit(String);
impl tracing::field::Visit for Visit {
fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
if field.name() == "message" {
self.0 = format!("{value:?}");
}
}
}
if *event.metadata().level() != Level::WARN {
return;
}
let sink = {
let guard = ACTIVE_SINK
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
match guard.as_ref() {
Some(sink) => Arc::clone(sink),
None => return,
}
};
let mut v = Visit(String::new());
event.record(&mut v);
sink.lock().expect("warn log").push(v.0);
}
}
struct HangingExecutor;
impl AgentExecutor for HangingExecutor {
fn execute<'a>(
&'a self,
_ctx: &'a RequestContext,
_queue: &'a dyn EventQueueWriter,
) -> Pin<Box<dyn Future<Output = A2aResult<()>> + Send + 'a>> {
Box::pin(async { Ok(()) })
}
fn on_shutdown<'a>(&'a self) -> Pin<Box<dyn Future<Output = ()> + Send + 'a>> {
Box::pin(async { std::future::pending::<()>().await })
}
}
static CAPTURE_LOCK: Mutex<()> = Mutex::new(());
static INSTALLED: OnceLock<()> = OnceLock::new();
fn warnings_during<F>(f: F) -> Vec<String>
where
F: FnOnce(),
{
let _serial = CAPTURE_LOCK
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
INSTALLED.get_or_init(|| {
let subscriber = Registry::default().with(WarnCapture);
let _ = tracing::subscriber::set_global_default(subscriber);
});
let sink: Sink = Arc::new(Mutex::new(Vec::new()));
*ACTIVE_SINK
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(Arc::clone(&sink));
f();
*ACTIVE_SINK
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) = None;
let out = sink.lock().expect("warn log").clone();
out
}
fn mentions_cleanup(warnings: &[String]) -> bool {
warnings.iter().any(|w| w.contains("executor cleanup"))
}
#[test]
fn clean_shutdown_warns_about_nothing() {
let rt = tokio::runtime::Builder::new_current_thread()
.enable_time()
.build()
.expect("runtime");
let warnings = warnings_during(|| {
rt.block_on(async {
let handler = make_handler();
let report = handler
.shutdown_with_timeout(Duration::from_millis(50))
.await;
assert!(
report.executor_cleanup_completed,
"the no-op executor's cleanup returns immediately"
);
});
});
assert!(
!mentions_cleanup(&warnings),
"a clean shutdown must not warn about executor cleanup; got {warnings:?}"
);
}
#[test]
fn fixed_budget_shutdown_warns_only_when_cleanup_hangs() {
let rt = || {
tokio::runtime::Builder::new_current_thread()
.enable_time()
.start_paused(true)
.build()
.expect("runtime")
};
let clean = warnings_during(|| {
rt().block_on(async {
let handler = make_handler();
let report = handler.shutdown().await;
assert!(report.executor_cleanup_completed);
});
});
assert!(
!mentions_cleanup(&clean),
"a clean shutdown() must not warn about executor cleanup; got {clean:?}"
);
let hung = warnings_during(|| {
rt().block_on(async {
let handler = RequestHandlerBuilder::new(HangingExecutor)
.build()
.expect("builder should succeed");
let report = handler.shutdown().await;
assert!(!report.executor_cleanup_completed);
});
});
assert!(
mentions_cleanup(&hung),
"a hung cleanup under shutdown() must be warned about; got {hung:?}"
);
}
#[test]
fn hung_cleanup_is_warned_about() {
let rt = tokio::runtime::Builder::new_current_thread()
.enable_time()
.build()
.expect("runtime");
let warnings = warnings_during(|| {
rt.block_on(async {
let handler = RequestHandlerBuilder::new(HangingExecutor)
.build()
.expect("builder should succeed");
let report = handler
.shutdown_with_timeout(Duration::from_millis(50))
.await;
assert!(
!report.executor_cleanup_completed,
"a hanging cleanup must be reported as incomplete"
);
});
});
assert!(
mentions_cleanup(&warnings),
"a hung executor cleanup must be warned about; got {warnings:?}"
);
}