1212)
1313
1414from ..spans import (
15- agent_workflow_span ,
1615 execute_tool_span ,
1716 update_execute_tool_span ,
1817 update_invoke_agent_span ,
2827from typing import TYPE_CHECKING , TypeVar
2928
3029if TYPE_CHECKING :
31- from typing import Any , AsyncIterator , Callable
30+ from typing import Any , Callable
3231
3332 from agents import Agent , Tool , ToolContext
3433
@@ -126,7 +125,6 @@ def _create_run_wrapper(
126125) -> "Callable[..., Any]" :
127126 """
128127 Wraps the agents.Runner.run methods to
129- - create and manage a root span for the agent workflow runs.
130128 - end the agent invocation span if an `AgentsException` is raised in `run()`.
131129
132130 Note agents.Runner.run_sync() is a wrapper around agents.Runner.run(),
@@ -150,88 +148,73 @@ async def wrapper(*args: "Any", **kwargs: "Any") -> "Any":
150148 else :
151149 agent = args [0 ].clone ()
152150
153- with agent_workflow_span (agent ) as workflow_span :
154- # Set conversation ID on workflow span early so it's captured even on errors
155- conversation_id = kwargs .get ("conversation_id" )
156- if conversation_id :
157- agent ._sentry_conversation_id = conversation_id
158-
159- workflow_span .set_attribute (
160- SPANDATA .GEN_AI_CONVERSATION_ID , conversation_id
161- )
162-
163- if "starting_agent" in kwargs :
164- kwargs ["starting_agent" ] = agent
165- else :
166- args = (agent , * args [1 :])
167-
168- try :
169- run_result = await original_func (* args , ** kwargs )
170- except AgentsException as exc :
171- exc_info = sys .exc_info ()
172- with capture_internal_exceptions ():
173- _capture_exception (exc )
174-
175- context_wrapper = getattr (exc .run_data , "context_wrapper" , None )
176- if context_wrapper is not None :
177- invoke_agent_span = getattr (
178- context_wrapper , "_sentry_agent_span" , None
151+ # Set conversation ID on workflow span early so it's captured even on errors
152+ conversation_id = kwargs .get ("conversation_id" )
153+ if conversation_id :
154+ agent ._sentry_conversation_id = conversation_id
155+
156+ if "starting_agent" in kwargs :
157+ kwargs ["starting_agent" ] = agent
158+ else :
159+ args = (agent , * args [1 :])
160+
161+ try :
162+ run_result = await original_func (* args , ** kwargs )
163+ except AgentsException as exc :
164+ exc_info = sys .exc_info ()
165+ with capture_internal_exceptions ():
166+ _capture_exception (exc )
167+
168+ context_wrapper = getattr (exc .run_data , "context_wrapper" , None )
169+ if context_wrapper is not None :
170+ invoke_agent_span = getattr (
171+ context_wrapper , "_sentry_agent_span" , None
172+ )
173+
174+ if (
175+ invoke_agent_span is not None
176+ and invoke_agent_span .end_timestamp is None
177+ ):
178+ update_invoke_agent_span (
179+ span = invoke_agent_span ,
180+ agent = agent ,
179181 )
180182
181- if (
182- invoke_agent_span is not None
183- and invoke_agent_span .end_timestamp is None
184- ):
185- update_invoke_agent_span (
186- span = invoke_agent_span ,
187- agent = agent ,
188- )
189-
190- invoke_agent_span .__exit__ (* exc_info )
191- delattr (context_wrapper , "_sentry_agent_span" )
192- reraise (* exc_info )
193- except Exception as exc :
194- exc_info = sys .exc_info ()
195- with capture_internal_exceptions ():
196- # Invoke agent span is not finished in this case.
197- # This is much less likely to occur than other cases because
198- # AgentRunner.run() is "just" a while loop around _run_single_turn.
199- _capture_exception (exc )
200- reraise (* exc_info )
201-
202- invoke_agent_span = getattr (
203- run_result .context_wrapper , "_sentry_agent_span" , None
204- )
205- if not invoke_agent_span :
206- return run_result
207-
208- update_invoke_agent_span (
209- span = invoke_agent_span ,
210- agent = agent ,
211- )
212-
213- invoke_agent_span .__exit__ (None , None , None )
214- delattr (run_result .context_wrapper , "_sentry_agent_span" )
183+ invoke_agent_span .__exit__ (* exc_info )
184+ delattr (context_wrapper , "_sentry_agent_span" )
185+ reraise (* exc_info )
186+ except Exception as exc :
187+ exc_info = sys .exc_info ()
188+ with capture_internal_exceptions ():
189+ # Invoke agent span is not finished in this case.
190+ # This is much less likely to occur than other cases because
191+ # AgentRunner.run() is "just" a while loop around _run_single_turn.
192+ _capture_exception (exc )
193+ reraise (* exc_info )
194+
195+ invoke_agent_span = getattr (
196+ run_result .context_wrapper , "_sentry_agent_span" , None
197+ )
198+ if not invoke_agent_span :
215199 return run_result
216200
201+ update_invoke_agent_span (
202+ span = invoke_agent_span ,
203+ agent = agent ,
204+ )
205+
206+ invoke_agent_span .__exit__ (None , None , None )
207+ delattr (run_result .context_wrapper , "_sentry_agent_span" )
208+ return run_result
209+
217210 return wrapper
218211
219212
220213def _create_run_streamed_wrapper (
221214 original_func : "Callable[..., Any]" ,
222215) -> "Callable[..., Any]" :
223216 """
224- Wraps the agents.Runner.run_streamed method to
225- - create a root span for streaming agent workflow runs.
226- - end the workflow span if and only if the response stream is consumed or cancelled.
227-
228- Unlike run(), run_streamed() returns immediately with a RunResultStreaming object
229- while execution continues in a background task. The workflow span must stay open
230- throughout the streaming operation and close when streaming completes or is abandoned.
231-
232- Note: We don't use isolation_scope() here because it uses context variables that
233- cannot span async boundaries (the __enter__ and __exit__ would be called from
234- different async contexts, causing ValueError).
217+ Wraps the agents.Runner.run_streamed method to inject run hooks.
235218 """
236219
237220 @wraps (original_func )
@@ -247,19 +230,6 @@ def wrapper(*args: "Any", **kwargs: "Any") -> "Any":
247230 if conversation_id :
248231 agent ._sentry_conversation_id = conversation_id
249232
250- # Start workflow span immediately (before run_streamed returns)
251- workflow_span = agent_workflow_span (agent )
252- workflow_span .__enter__ ()
253-
254- # Set conversation ID on workflow span early so it's captured even on errors
255- if conversation_id :
256- workflow_span .set_attribute (
257- SPANDATA .GEN_AI_CONVERSATION_ID , conversation_id
258- )
259-
260- # Store span on agent for cleanup
261- agent ._sentry_workflow_span = workflow_span
262-
263233 if "starting_agent" in kwargs :
264234 kwargs ["starting_agent" ] = agent
265235 else :
@@ -276,45 +246,9 @@ def wrapper(*args: "Any", **kwargs: "Any") -> "Any":
276246 # Call original function to get RunResultStreaming
277247 run_result = original_func (* args , ** kwargs )
278248 except Exception as exc :
279- # If run_streamed itself fails (not the background task), clean up immediately
280- workflow_span .__exit__ (* sys .exc_info ())
281249 _capture_exception (exc )
282250 raise
283251
284- def _close_workflow_span () -> None :
285- if hasattr (agent , "_sentry_workflow_span" ):
286- workflow_span .__exit__ (* sys .exc_info ())
287- delattr (agent , "_sentry_workflow_span" )
288-
289- if hasattr (run_result , "stream_events" ):
290- original_stream_events = run_result .stream_events
291-
292- @wraps (original_stream_events )
293- async def wrapped_stream_events (
294- * stream_args : "Any" , ** stream_kwargs : "Any"
295- ) -> "AsyncIterator[Any]" :
296- try :
297- async for event in original_stream_events (
298- * stream_args , ** stream_kwargs
299- ):
300- yield event
301- finally :
302- _close_workflow_span ()
303-
304- run_result .stream_events = wrapped_stream_events
305-
306- if hasattr (run_result , "cancel" ):
307- original_cancel = run_result .cancel
308-
309- @wraps (original_cancel )
310- def wrapped_cancel (* cancel_args : "Any" , ** cancel_kwargs : "Any" ) -> "Any" :
311- try :
312- return original_cancel (* cancel_args , ** cancel_kwargs )
313- finally :
314- _close_workflow_span ()
315-
316- run_result .cancel = wrapped_cancel
317-
318252 return run_result
319253
320254 return wrapper
0 commit comments