@@ -91,6 +91,34 @@ async def _detach_tracing(detach: Any, tracer: Tracer) -> None:
9191 await tracer .shutdown ()
9292
9393
94+ async def test_suspended_outcome_without_snapshot_emits_no_event () -> None :
95+ agent , _channel , _exporter , tracer = _build_suspending_agent ("thread-no-snapshot" )
96+ observed : list [str ] = []
97+
98+ async def observe (event , signal = None ): # noqa: ANN001, ANN202
99+ del signal
100+ observed .append (event .type )
101+
102+ unsubscribe = agent .subscribe (observe )
103+ await agent ._emit_suspended_event ("suspended" )
104+ unsubscribe ()
105+ await tracer .shutdown ()
106+
107+ assert observed == []
108+
109+
110+ async def test_recorder_suspension_without_active_run_is_noop () -> None :
111+ from cubepi .tracing .recorder import Recorder
112+
113+ tracer = Tracer (service_name = "test" , exporters = [])
114+ recorder = Recorder (tracer )
115+
116+ recorder ._on_agent_suspended ()
117+ await tracer .shutdown ()
118+
119+ assert recorder ._run is None
120+
121+
94122async def test_suspended_run_exports_suspended_trace_without_abort_or_leaks () -> None :
95123 agent , channel , exporter , tracer = _build_suspending_agent ("thread-1" )
96124 provider_stack_baseline = len (mcp_tracing ._provider_stack )
@@ -130,6 +158,105 @@ async def test_suspended_run_exports_suspended_trace_without_abort_or_leaks() ->
130158 assert len (mcp_tracing ._active_entries ) == active_entries_baseline
131159
132160
161+ async def test_recorder_suspension_finalizes_all_open_resources () -> None :
162+ from cubepi .tracing .recorder import Recorder , _RunState
163+
164+ class _Span :
165+ def __init__ (self ) -> None :
166+ self .attributes : dict [str , Any ] = {}
167+ self .ended = False
168+
169+ def set_attribute (self , key : str , value : Any ) -> None :
170+ self .attributes [key ] = value
171+
172+ def end (self ) -> None :
173+ self .ended = True
174+
175+ class _Stream :
176+ def __init__ (self ) -> None :
177+ self .closed = False
178+
179+ def close (self ) -> None :
180+ self .closed = True
181+
182+ tracer = Tracer (service_name = "test" , exporters = [])
183+ recorder = Recorder (tracer )
184+ agent_span = _Span ()
185+ turn_span = _Span ()
186+ chat_span = _Span ()
187+ tool_span = _Span ()
188+ stream = _Stream ()
189+ run = _RunState (
190+ run_id = "run-open-resources" ,
191+ agent_span = agent_span ,
192+ turn_span = turn_span ,
193+ chat_span = chat_span ,
194+ tool_spans = {"tool-1" : tool_span },
195+ stream_file = stream ,
196+ )
197+ recorder ._run = run
198+
199+ recorder ._on_agent_suspended ()
200+ await tracer .shutdown ()
201+
202+ for span in (agent_span , turn_span , chat_span , tool_span ):
203+ assert span .attributes ["cubepi.run.outcome" ] == "suspended"
204+ assert span .ended is True
205+ assert stream .closed is True
206+ assert run .tool_spans == {}
207+ assert run .chat_span is None
208+ assert run .turn_span is None
209+ assert run .stream_file is None
210+ assert recorder ._run is None
211+
212+
213+ async def test_recorder_suspension_cleanup_continues_after_resource_errors () -> None :
214+ from cubepi .tracing .recorder import Recorder , _RunState
215+
216+ class _Span :
217+ def __init__ (self ) -> None :
218+ self .attributes : dict [str , Any ] = {}
219+ self .ended = False
220+
221+ def set_attribute (self , key : str , value : Any ) -> None :
222+ self .attributes [key ] = value
223+
224+ def end (self ) -> None :
225+ self .ended = True
226+
227+ class _BoomSpan (_Span ):
228+ def end (self ) -> None :
229+ raise RuntimeError ("span end failed" )
230+
231+ class _BoomStream :
232+ def close (self ) -> None :
233+ raise OSError ("stream close failed" )
234+
235+ tracer = Tracer (service_name = "test" , exporters = [])
236+ recorder = Recorder (tracer )
237+ agent_span = _Span ()
238+ run = _RunState (
239+ run_id = "run-resource-errors" ,
240+ agent_span = agent_span ,
241+ turn_span = _BoomSpan (),
242+ chat_span = _BoomSpan (),
243+ tool_spans = {"tool-1" : _BoomSpan ()},
244+ stream_file = _BoomStream (),
245+ )
246+ recorder ._run = run
247+
248+ recorder ._on_agent_suspended ()
249+ await tracer .shutdown ()
250+
251+ assert agent_span .attributes ["cubepi.run.outcome" ] == "suspended"
252+ assert agent_span .ended is True
253+ assert run .tool_spans == {}
254+ assert run .chat_span is None
255+ assert run .turn_span is None
256+ assert run .stream_file is None
257+ assert recorder ._run is None
258+
259+
133260async def test_suspension_event_observes_committed_runtime_state () -> None :
134261 agent , channel , _exporter , tracer = _build_suspending_agent ("thread-2" )
135262 detach_tracing = tracer .attach (agent )
0 commit comments