@@ -146,7 +146,7 @@ where
146146/// Isolated in-memory span exporters for tracing tests.
147147#[ cfg( test) ]
148148pub mod test_exporter {
149- use std:: sync:: OnceLock ;
149+ use std:: sync:: { Arc , OnceLock } ;
150150
151151 use tracing:: { Dispatch , Subscriber , dispatcher:: WeakDispatch } ;
152152
@@ -159,6 +159,7 @@ pub mod test_exporter {
159159 struct CloseWithDispatch < S > {
160160 inner : S ,
161161 dispatch : OnceLock < WeakDispatch > ,
162+ closed : Arc < tokio:: sync:: Notify > ,
162163 }
163164
164165 impl < S : Subscriber > Subscriber for CloseWithDispatch < S > {
@@ -224,7 +225,14 @@ pub mod test_exporter {
224225 . get ( )
225226 . and_then ( WeakDispatch :: upgrade)
226227 . expect ( "a live span keeps its test dispatcher alive" ) ;
227- tracing:: dispatcher:: with_default ( & dispatch, || self . inner . try_close ( id) )
228+ let closed = tracing:: dispatcher:: with_default ( & dispatch, || self . inner . try_close ( id) ) ;
229+ if closed {
230+ // The simple exporter has finished before try_close returns.
231+ // Wake assertions only after the owning registry and layers
232+ // have completed cleanup, including recursive parent closure.
233+ self . closed . notify_waiters ( ) ;
234+ }
235+ closed
228236 }
229237
230238 fn current_span ( & self ) -> tracing_core:: span:: Current {
@@ -273,14 +281,17 @@ pub mod test_exporter {
273281 . with_simple_exporter ( exporter. clone ( ) )
274282 . build ( ) ;
275283 let subscriber = tracing_subscriber:: registry ( ) . with ( super :: layer ( & provider, None ) ) ;
284+ let closed = Arc :: new ( tokio:: sync:: Notify :: new ( ) ) ;
276285 let dispatch = Dispatch :: new ( CloseWithDispatch {
277286 inner : subscriber,
278287 dispatch : OnceLock :: new ( ) ,
288+ closed : Arc :: clone ( & closed) ,
279289 } ) ;
280290 TracingTestGuard {
281291 _default : tracing:: dispatcher:: set_default ( & dispatch) ,
282292 _provider : provider,
283293 exporter,
294+ closed,
284295 _lock : lock,
285296 }
286297 }
@@ -291,6 +302,52 @@ pub mod test_exporter {
291302 self . exporter . get_finished_spans ( ) . expect ( "in-memory spans" )
292303 }
293304
305+ /// Wait for expected spans to finish before taking an assertion snapshot.
306+ ///
307+ /// A completed `SQLx` query may still have its span held by a `SQLite`
308+ /// worker. Export is synchronous once the span closes, but flushing
309+ /// cannot close that live span. Await closure notifications instead of
310+ /// assuming the query result also means tracing cleanup has completed.
311+ pub async fn wait_for_spans (
312+ & self ,
313+ predicate : impl Fn ( & [ opentelemetry_sdk:: trace:: SpanData ] ) -> bool + Send + Sync ,
314+ ) -> Vec < opentelemetry_sdk:: trace:: SpanData > {
315+ tokio:: time:: timeout ( std:: time:: Duration :: from_secs ( 5 ) , async {
316+ loop {
317+ // notify_waiters wakes futures created before notification,
318+ // even before polling. Subscribe before reading so closure
319+ // between the snapshot and await cannot lose a wakeup.
320+ let notified = self . closed . notified ( ) ;
321+ let spans = self . finished_spans ( ) ;
322+ if predicate ( & spans) {
323+ return spans;
324+ }
325+ notified. await ;
326+ }
327+ } )
328+ . await
329+ . unwrap_or_else ( |_| {
330+ panic ! (
331+ "timed out waiting for expected spans, got {:?}" ,
332+ self . finished_spans( )
333+ . iter( )
334+ . map( |span| & span. name)
335+ . collect:: <Vec <_>>( )
336+ )
337+ } )
338+ }
339+
340+ /// Wait for the completed span named `name`.
341+ pub async fn wait_for_span ( & self , name : & str ) -> opentelemetry_sdk:: trace:: SpanData {
342+ let spans = self
343+ . wait_for_spans ( |spans| spans. iter ( ) . any ( |span| span. name == name) )
344+ . await ;
345+ spans
346+ . into_iter ( )
347+ . find ( |span| span. name == name)
348+ . expect ( "the awaited snapshot contains the expected span" )
349+ }
350+
294351 /// Spans named `name`.
295352 pub fn spans_named ( & self , name : & str ) -> Vec < opentelemetry_sdk:: trace:: SpanData > {
296353 self . finished_spans ( )
@@ -378,6 +435,7 @@ pub mod test_exporter {
378435 _default : tracing:: dispatcher:: DefaultGuard ,
379436 _provider : opentelemetry_sdk:: trace:: SdkTracerProvider ,
380437 exporter : opentelemetry_sdk:: trace:: InMemorySpanExporter ,
438+ closed : Arc < tokio:: sync:: Notify > ,
381439 _lock : std:: sync:: MutexGuard < ' static , ( ) > ,
382440 }
383441
@@ -690,6 +748,56 @@ mod tests {
690748 ) ;
691749 }
692750
751+ #[ tokio:: test]
752+ async fn tracing_waits_for_worker_held_spans_to_close ( ) {
753+ let traced = test_exporter:: install_traced ( ) ;
754+ let parent = tracing:: info_span!( "awaited_parent" ) ;
755+ let child = tracing:: info_span!( parent: & parent, "awaited_child" ) ;
756+ drop ( parent) ;
757+
758+ let ( release, released) = std:: sync:: mpsc:: channel ( ) ;
759+ let worker = std:: thread:: spawn ( move || {
760+ released. recv ( ) . expect ( "test releases the worker's span" ) ;
761+ drop ( child) ;
762+ } ) ;
763+
764+ assert ! ( traced. finished_spans( ) . is_empty( ) ) ;
765+ let waiting = traced. wait_for_span ( "awaited_parent" ) ;
766+ tokio:: pin!( waiting) ;
767+ assert ! ( futures:: poll!( waiting. as_mut( ) ) . is_pending( ) ) ;
768+
769+ release. send ( ( ) ) . unwrap ( ) ;
770+ let parent = waiting. await ;
771+ worker. join ( ) . expect ( "worker closes child and parent" ) ;
772+ let child = traced. span_named ( "awaited_child" ) ;
773+ test_exporter:: assert_is_root ( & parent) ;
774+ assert_eq ! ( child. parent_span_id, parent. span_context. span_id( ) ) ;
775+ }
776+
777+ #[ tokio:: test]
778+ async fn tracing_wait_does_not_lose_closure_between_snapshot_and_await ( ) {
779+ let traced = test_exporter:: install_traced ( ) ;
780+ let span = std:: sync:: Mutex :: new ( Some ( tracing:: info_span!( "close_before_await" ) ) ) ;
781+ let spans = traced
782+ . wait_for_spans ( |spans| {
783+ // The first snapshot is empty. Close its span before polling
784+ // the notification future, as a worker could do concurrently.
785+ drop ( span. lock ( ) . unwrap ( ) . take ( ) ) ;
786+ spans. iter ( ) . any ( |span| span. name == "close_before_await" )
787+ } )
788+ . await ;
789+ assert_eq ! ( spans. len( ) , 1 ) ;
790+ test_exporter:: assert_is_root ( & spans[ 0 ] ) ;
791+ }
792+
793+ #[ tokio:: test( start_paused = true ) ]
794+ #[ should_panic( expected = "timed out waiting for expected spans" ) ]
795+ async fn tracing_wait_times_out_when_a_span_never_closes ( ) {
796+ let traced = test_exporter:: install_traced ( ) ;
797+ let _span = tracing:: info_span!( "still_open" ) ;
798+ traced. wait_for_span ( "still_open" ) . await ;
799+ }
800+
693801 #[ tokio:: test]
694802 async fn tracing_events_are_not_exported ( ) {
695803 let traced = test_exporter:: install_traced ( ) ;
0 commit comments