@@ -120,6 +120,43 @@ async def _ticker() -> None:
120120 await tracer .shutdown ()
121121
122122
123+ async def test_meter_force_flush_calls_provider_off_the_loop_thread ():
124+ # Meter mirrors Tracer: its provider force_flush is synchronous and
125+ # must run in a worker thread, never on the loop.
126+ from opentelemetry .sdk .metrics import MeterProvider
127+ from opentelemetry .sdk .metrics .export import InMemoryMetricReader
128+ from opentelemetry .sdk .resources import Resource
129+
130+ from cubepi .tracing import Meter
131+ from cubepi .tracing .schema import SCHEMA_URL
132+
133+ reader = InMemoryMetricReader ()
134+ resource = Resource .create ({"service.name" : "test" }, schema_url = SCHEMA_URL )
135+ provider = MeterProvider (resource = resource , metric_readers = [reader ])
136+ meter = Meter .__new__ (Meter )
137+ meter ._provider = provider # type: ignore[attr-defined]
138+ meter ._shutdown = False # type: ignore[attr-defined]
139+
140+ loop_thread = threading .current_thread ().name
141+ flush_threads : list [str ] = []
142+ inner = provider .force_flush
143+
144+ def _spy (timeout_millis : int = 30_000 ) -> bool :
145+ flush_threads .append (threading .current_thread ().name )
146+ return inner (timeout_millis = timeout_millis )
147+
148+ provider .force_flush = _spy # type: ignore[method-assign]
149+ try :
150+ assert await meter .force_flush () is True
151+ assert flush_threads , "provider.force_flush must have run"
152+ assert flush_threads [0 ] != loop_thread , (
153+ "sync metric flush must run in a worker thread, not on the loop"
154+ )
155+ finally :
156+ provider .force_flush = inner # type: ignore[method-assign]
157+ provider .shutdown ()
158+
159+
123160async def test_trace_background_exits_before_export_completes ():
124161 exporter = _SlowExporter (export_seconds = 0.5 )
125162 agent , tracer = _build (exporter )
0 commit comments