@@ -240,42 +240,28 @@ async def stream_messages() -> None:
240240 task_id = task .id ,
241241 timeout = 90 , # Increased timeout for CI environments
242242 ):
243+ # A turn emits several messages (user echo, reasoning, agent text),
244+ # each ending in "full" or "done"; consume until the text reply lands.
243245 msg_type = event .get ("type" )
244246 if msg_type == "full" :
245- task_message_update = StreamTaskMessageFull .model_validate (event )
246- if task_message_update .parent_task_message and task_message_update .parent_task_message .id :
247- finished_message = await client .messages .retrieve (task_message_update .parent_task_message .id )
248- if (
249- finished_message .content
250- and finished_message .content .type == "text"
251- and finished_message .content .author == "user"
252- ):
253- user_message_found = True
254- elif (
255- finished_message .content
256- and finished_message .content .type == "text"
257- and finished_message .content .author == "agent"
258- ):
259- agent_response_found = True
260- elif finished_message .content and finished_message .content .type == "reasoning" :
261- reasoning_found = True
262-
263- # Exit early if we have what we need
264- if user_message_found and agent_response_found :
265- break
266-
247+ parent_task_message = StreamTaskMessageFull .model_validate (event ).parent_task_message
267248 elif msg_type == "done" :
268- task_message_update_done = StreamTaskMessageDone .model_validate (event )
269- if task_message_update_done .parent_task_message and task_message_update_done .parent_task_message .id :
270- finished_message = await client .messages .retrieve (task_message_update_done .parent_task_message .id )
271- if finished_message .content and finished_message .content .type == "reasoning" :
272- reasoning_found = True
273- elif (
274- finished_message .content
275- and finished_message .content .type == "text"
276- and finished_message .content .author == "agent"
277- ):
278- agent_response_found = True
249+ parent_task_message = StreamTaskMessageDone .model_validate (event ).parent_task_message
250+ else :
251+ continue
252+
253+ if parent_task_message and parent_task_message .id :
254+ finished_message = await client .messages .retrieve (parent_task_message .id )
255+ content = finished_message .content
256+ if content and content .type == "text" and content .author == "user" :
257+ user_message_found = True
258+ elif content and content .type == "text" and content .author == "agent" :
259+ agent_response_found = True
260+ elif content and content .type == "reasoning" :
261+ reasoning_found = True
262+
263+ # Stop once both the user echo and the agent's text reply are seen.
264+ if user_message_found and agent_response_found :
279265 break
280266
281267 stream_task = asyncio .create_task (stream_messages ())
0 commit comments