@@ -206,24 +206,24 @@ def _build_context_metadata(event: Event, ctx: InvocationContext) -> Dict[str, A
206206 return metadata
207207
208208
209- def _build_message_metadata (event : Event ) -> Dict [str , Any ]:
209+ def _build_message_metadata (event : Event , effective_id : str ) -> Dict [str , Any ]:
210210 """Build message/event metadata (object_type, tag, llm_response_id)."""
211211 return {
212212 MESSAGE_METADATA_OBJECT_TYPE_KEY : _infer_message_object_type (event ) or "" ,
213213 MESSAGE_METADATA_TAG_KEY : _infer_message_tag (event ),
214- MESSAGE_METADATA_RESPONSE_ID_KEY : event . response_id or "" ,
214+ MESSAGE_METADATA_RESPONSE_ID_KEY : effective_id ,
215215 }
216216
217217
218- def _build_event_metadata (event : Event , message : Message , ctx : InvocationContext ) -> Dict [str , Any ]:
218+ def _build_event_metadata (event : Event , message : Message , ctx : InvocationContext , effective_id : str ) -> Dict [str , Any ]:
219219 metadata = _build_context_metadata (event , ctx )
220- msg_meta = _build_message_metadata (event )
220+ msg_meta = _build_message_metadata (event , effective_id )
221221 set_metadata (metadata , MESSAGE_METADATA_OBJECT_TYPE_KEY , msg_meta .get (MESSAGE_METADATA_OBJECT_TYPE_KEY ) or "" )
222222 set_metadata (metadata , MESSAGE_METADATA_TAG_KEY , msg_meta .get (MESSAGE_METADATA_TAG_KEY ) or "" )
223223 set_metadata (metadata , MESSAGE_METADATA_RESPONSE_ID_KEY , msg_meta .get (MESSAGE_METADATA_RESPONSE_ID_KEY ) or "" )
224- if any (
225- get_metadata (p .root .metadata , A2A_DATA_PART_METADATA_TYPE_KEY ) ==
226- A2A_DATA_PART_METADATA_TYPE_STREAMING_FUNCTION_CALL_DELTA for p in message .parts if p .root .metadata ):
224+ streaming_delta = A2A_DATA_PART_METADATA_TYPE_STREAMING_FUNCTION_CALL_DELTA
225+ if any ( get_metadata (p .root .metadata , A2A_DATA_PART_METADATA_TYPE_KEY ) == streaming_delta
226+ for p in message .parts if p .root .metadata ):
227227 set_metadata (metadata , "streaming_tool_call" , "true" )
228228 return metadata
229229
@@ -234,27 +234,41 @@ def _mark_long_running_tools(a2a_parts: List[A2APart], event: Event) -> None:
234234 return
235235 for a2a_part in a2a_parts :
236236 root = a2a_part .root
237- if (isinstance (root , DataPart ) and root .metadata and get_metadata (
238- root .metadata , A2A_DATA_PART_METADATA_TYPE_KEY ) == A2A_DATA_PART_METADATA_TYPE_FUNCTION_CALL
239- and root .data .get ("id" ) in event .long_running_tool_ids ):
240- set_metadata (root .metadata , A2A_DATA_PART_METADATA_IS_LONG_RUNNING_KEY , True )
237+ if not isinstance (root , DataPart ) or not root .metadata :
238+ continue
239+ if get_metadata (root .metadata , A2A_DATA_PART_METADATA_TYPE_KEY ) != A2A_DATA_PART_METADATA_TYPE_FUNCTION_CALL :
240+ continue
241+ if root .data .get ("id" ) not in event .long_running_tool_ids :
242+ continue
243+ set_metadata (root .metadata , A2A_DATA_PART_METADATA_IS_LONG_RUNNING_KEY , True )
244+
245+
246+ def _effective_response_id (event : Event ) -> str :
247+ """Return ``response_id`` when present, otherwise a new UUID.
248+
249+ Callers that need the same id across multiple locations should invoke this
250+ once and pass the result explicitly.
251+ """
252+ return event .response_id or str (uuid .uuid4 ())
241253
242254
243- def _build_message (event : Event , a2a_parts : List [A2APart ], role : Role ) -> Optional [Message ]:
255+ def _build_message (event : Event , a2a_parts : List [A2APart ], role : Role , effective_id : str ) -> Optional [Message ]:
244256 """Assemble an A2A Message from converted parts, or return None if empty."""
245257 if not a2a_parts :
246258 return None
247- message_id = event .response_id or str (uuid .uuid4 ())
248- message = Message (message_id = message_id , role = role , parts = a2a_parts )
249- msg_meta = _build_message_metadata (event )
259+ message = Message (message_id = effective_id , role = role , parts = a2a_parts )
260+ msg_meta = _build_message_metadata (event , effective_id )
250261 if msg_meta :
251262 message .metadata = msg_meta
252263 return message
253264
254265
255266def _is_streaming_delta (a2a_part : A2APart ) -> bool :
256- return (a2a_part .root .metadata is not None and get_metadata (a2a_part .root .metadata , A2A_DATA_PART_METADATA_TYPE_KEY )
257- == A2A_DATA_PART_METADATA_TYPE_STREAMING_FUNCTION_CALL_DELTA )
267+ meta = a2a_part .root .metadata
268+ if meta is None :
269+ return False
270+ t = get_metadata (meta , A2A_DATA_PART_METADATA_TYPE_KEY )
271+ return t == A2A_DATA_PART_METADATA_TYPE_STREAMING_FUNCTION_CALL_DELTA
258272
259273
260274def _collect_parts (
@@ -320,7 +334,8 @@ def convert_event_to_a2a_message(
320334 return None
321335
322336 a2a_parts = _collect_parts (event , ** rules )
323- return _build_message (event , a2a_parts , role )
337+ effective_id = _effective_response_id (event )
338+ return _build_message (event , a2a_parts , role , effective_id )
324339
325340
326341def convert_content_to_a2a_message (
@@ -425,8 +440,8 @@ def convert_a2a_message_to_event(
425440 if gpart is None :
426441 logger .warning ("Failed to convert A2A part, skipping: %s" , a2a_part )
427442 continue
428- if ( metadata_is_true (a2a_part .root .metadata , A2A_DATA_PART_METADATA_IS_LONG_RUNNING_KEY )
429- and gpart .function_call ) :
443+ is_lr = metadata_is_true (a2a_part .root .metadata , A2A_DATA_PART_METADATA_IS_LONG_RUNNING_KEY )
444+ if is_lr and gpart .function_call :
430445 long_running_tool_ids .add (gpart .function_call .id )
431446 parts .append (gpart )
432447 except Exception as ex : # pylint: disable=broad-except
@@ -436,8 +451,8 @@ def convert_a2a_message_to_event(
436451 if not parts :
437452 logger .warning ("No parts could be converted from A2A message %s" , a2a_message )
438453
439- object_type = ( get_metadata (msg_meta , MESSAGE_METADATA_OBJECT_TYPE_KEY )
440- or _infer_a2a_message_object_type (parts , partial = partial ) or _default_object_type (partial ) )
454+ ot = get_metadata (msg_meta , MESSAGE_METADATA_OBJECT_TYPE_KEY )
455+ object_type = ot or _infer_a2a_message_object_type (parts , partial = partial ) or _default_object_type (partial )
441456
442457 return Event (
443458 invocation_id = inv_id ,
@@ -590,31 +605,55 @@ def _create_error_status_event(
590605 )
591606
592607
608+ def _a2a_part_requests_euc_auth (part : A2APart ) -> bool :
609+ root = part .root
610+ md = root .metadata
611+ if not md :
612+ return False
613+ t = get_metadata (md , A2A_DATA_PART_METADATA_TYPE_KEY )
614+ return all (
615+ [
616+ t == A2A_DATA_PART_METADATA_TYPE_FUNCTION_CALL ,
617+ metadata_is_true (md , A2A_DATA_PART_METADATA_IS_LONG_RUNNING_KEY ),
618+ root .data .get ("name" ) == REQUEST_EUC_FUNCTION_CALL_NAME ,
619+ ]
620+ )
621+
622+
623+ def _a2a_part_is_long_running_function_call (part : A2APart ) -> bool :
624+ root = part .root
625+ md = root .metadata
626+ if not md :
627+ return False
628+ t = get_metadata (md , A2A_DATA_PART_METADATA_TYPE_KEY )
629+ return all (
630+ [
631+ t == A2A_DATA_PART_METADATA_TYPE_FUNCTION_CALL ,
632+ metadata_is_true (md , A2A_DATA_PART_METADATA_IS_LONG_RUNNING_KEY ),
633+ ]
634+ )
635+
636+
593637def _create_status_update_event (
594638 message : Message ,
595639 ctx : InvocationContext ,
596640 event : Event ,
597641 task_id : Optional [str ],
598642 context_id : Optional [str ],
643+ effective_id : str = "" ,
599644) -> TaskStatusUpdateEvent :
600645 status = TaskStatus (state = TaskState .working , message = message , timestamp = _now_iso ())
601646
602- if any (
603- get_metadata (p .root .metadata , A2A_DATA_PART_METADATA_TYPE_KEY ) == A2A_DATA_PART_METADATA_TYPE_FUNCTION_CALL
604- and metadata_is_true (p .root .metadata , A2A_DATA_PART_METADATA_IS_LONG_RUNNING_KEY )
605- and p .root .data .get ("name" ) == REQUEST_EUC_FUNCTION_CALL_NAME for p in message .parts if p .root .metadata ):
647+ if any (_a2a_part_requests_euc_auth (p ) for p in message .parts ):
606648 status .state = TaskState .auth_required
607- elif any (
608- get_metadata (p .root .metadata , A2A_DATA_PART_METADATA_TYPE_KEY ) == A2A_DATA_PART_METADATA_TYPE_FUNCTION_CALL
609- and metadata_is_true (p .root .metadata , A2A_DATA_PART_METADATA_IS_LONG_RUNNING_KEY ) for p in message .parts
610- if p .root .metadata ):
649+ elif any (_a2a_part_is_long_running_function_call (p ) for p in message .parts ):
611650 status .state = TaskState .input_required
612651
613652 return TaskStatusUpdateEvent (
614653 task_id = task_id ,
615654 context_id = context_id ,
616655 status = status ,
617- metadata = _build_event_metadata (event , message , ctx ),
656+ metadata = _build_event_metadata (event , message , ctx , effective_id ),
618657 final = False ,
619658 )
620659
@@ -626,8 +665,9 @@ def _create_artifact_update_event(
626665 task_id : Optional [str ] = None ,
627666 context_id : Optional [str ] = None ,
628667 last_chunk : bool = False ,
668+ effective_id : str = "" ,
629669) -> TaskArtifactUpdateEvent :
630- artifact_id = "" if last_chunk else ( event . response_id or "" )
670+ artifact_id = "" if last_chunk else effective_id
631671 return TaskArtifactUpdateEvent (
632672 task_id = task_id ,
633673 context_id = context_id ,
@@ -636,7 +676,7 @@ def _create_artifact_update_event(
636676 parts = [] if last_chunk else message .parts ,
637677 ),
638678 last_chunk = last_chunk ,
639- metadata = _build_event_metadata (event , message , ctx ),
679+ metadata = _build_event_metadata (event , message , ctx , effective_id ),
640680 )
641681
642682
@@ -674,12 +714,14 @@ def _notify(evt: A2AEvent) -> None:
674714
675715 message = convert_event_to_a2a_message (event , invocation_context )
676716 if message :
717+ effective_id = message .message_id
677718 status_event = _create_status_update_event (
678719 message ,
679720 invocation_context ,
680721 event ,
681722 task_id ,
682723 context_id ,
724+ effective_id = effective_id ,
683725 )
684726 _notify (status_event )
685727
@@ -691,6 +733,7 @@ def _notify(evt: A2AEvent) -> None:
691733 task_id = task_id ,
692734 context_id = context_id ,
693735 last_chunk = False ,
736+ effective_id = effective_id ,
694737 )
695738 a2a_events .append (artifact_event )
696739
0 commit comments