@@ -20,20 +20,22 @@ async def execute(
2020 event_queue : EventQueue ,
2121 input_queue : AgentInputQueue ,
2222 ) -> None :
23+ turn = await input_queue .get ()
24+
2325 if (
24- (not context .message )
25- or (not context .message .task_id )
26- or (not context .message .context_id )
26+ (not turn .message )
27+ or (not turn .message .task_id )
28+ or (not turn .message .context_id )
2729 ):
2830 raise Exception ("invalid message" )
2931
30- logging .debug (f"received message: { context .message } " )
32+ logging .debug (f"received message: { turn .message } " )
3133
3234 # The V2 request handler requires an initial Task to be enqueued
3335 # before any status/artifact update events are emitted.
34- task = context .current_task
36+ task = turn .current_task
3537 if task is None :
36- task = new_task_from_user_message (context .message )
38+ task = new_task_from_user_message (turn .message )
3739 await event_queue .enqueue_event (task )
3840
3941 task_updater = TaskUpdater (
@@ -42,14 +44,14 @@ async def execute(
4244 context_id = task .context_id ,
4345 )
4446
45- if context .message .parts [0 ].WhichOneof ("content" ) != "text" :
47+ if turn .message .parts [0 ].WhichOneof ("content" ) != "text" :
4648 raise Exception ("only text parts are supported" )
4749
48- result = await self .agent .invoke (context .message .parts [0 ].text )
50+ result = await self .agent .invoke (turn .message .parts [0 ].text )
4951
5052 response = Message (
5153 role = Role .ROLE_AGENT ,
52- message_id = context .message .message_id ,
54+ message_id = turn .message .message_id ,
5355 parts = [Part (text = result )],
5456 )
5557 await task_updater .add_artifact (
0 commit comments