Skip to content

Commit 6896769

Browse files
authored
Merge branch 'main' into codex/artifact-ref-scope
2 parents a3a91ba + 2e28e5d commit 6896769

8 files changed

Lines changed: 237 additions & 7 deletions

File tree

.github/workflows/analyze-releases-for-adk-docs-updates.yml

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -52,7 +52,9 @@ jobs:
5252
uses: actions/cache/restore@v4
5353
with:
5454
path: contributing/samples/adk_team/adk_documentation/adk_release_analyzer/sessions.db
55-
key: analyzer-session-db
55+
key: analyzer-session-db-${{ github.run_id }}-${{ github.run_attempt }}
56+
restore-keys: |
57+
analyzer-session-db-
5658
5759
- name: Run Analyzing Script
5860
env:
@@ -88,4 +90,4 @@ jobs:
8890
uses: actions/cache/save@v4
8991
with:
9092
path: contributing/samples/adk_team/adk_documentation/adk_release_analyzer/sessions.db
91-
key: analyzer-session-db
93+
key: analyzer-session-db-${{ github.run_id }}-${{ github.run_attempt }}

CONTRIBUTING.md

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -133,7 +133,7 @@ part before or alongside your code PR.
133133
1. **Clone the repository:**
134134

135135
```shell
136-
gh repo clone google/adk-python -- -b v2
136+
gh repo clone google/adk-python
137137
cd adk-python
138138
```
139139

src/google/adk/a2a/executor/a2a_agent_executor.py

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -251,8 +251,10 @@ async def _handle_request(
251251
)
252252

253253
task_result_aggregator = TaskResultAggregator()
254+
last_adk_event = None
254255
async with Aclosing(runner.run_async(**vars(run_request))) as agen:
255256
async for adk_event in agen:
257+
last_adk_event = adk_event
256258
for a2a_event in self._config.event_converter(
257259
adk_event,
258260
invocation_context,
@@ -270,6 +272,22 @@ async def _handle_request(
270272
task_result_aggregator.process_event(e)
271273
await event_queue.enqueue_event(e)
272274

275+
# Build metadata for final event to preserve invocation_id and event_id.
276+
final_metadata = {
277+
_get_adk_metadata_key('app_name'): runner.app_name,
278+
_get_adk_metadata_key('user_id'): run_request.user_id,
279+
_get_adk_metadata_key('session_id'): run_request.session_id,
280+
}
281+
if last_adk_event:
282+
for key, attr in [
283+
('invocation_id', 'invocation_id'),
284+
('author', 'author'),
285+
('event_id', 'id'),
286+
]:
287+
val = getattr(last_adk_event, attr, None)
288+
if val is not None:
289+
final_metadata[_get_adk_metadata_key(key)] = val
290+
273291
# publish the task result event - this is final
274292
if (
275293
task_result_aggregator.task_state == TaskState.working
@@ -287,6 +305,7 @@ async def _handle_request(
287305
artifact_id=platform_uuid.new_uuid(),
288306
parts=task_result_aggregator.task_status_message.parts,
289307
),
308+
metadata=final_metadata,
290309
)
291310
)
292311
# public the final status update event
@@ -299,6 +318,7 @@ async def _handle_request(
299318
).isoformat(),
300319
),
301320
context_id=context.context_id,
321+
metadata=final_metadata,
302322
final=True,
303323
)
304324
else:
@@ -312,6 +332,7 @@ async def _handle_request(
312332
message=task_result_aggregator.task_status_message,
313333
),
314334
context_id=context.context_id,
335+
metadata=final_metadata,
315336
final=True,
316337
)
317338

src/google/adk/flows/llm_flows/contents.py

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -108,7 +108,6 @@ def _rearrange_events_for_async_function_responses_in_history(
108108
events: list[Event],
109109
) -> list[Event]:
110110
"""Rearrange the async function_response events in the history."""
111-
112111
function_call_id_to_response_events_index: dict[str, int] = {}
113112
for i, event in enumerate(events):
114113
function_responses = event.get_function_responses()
@@ -117,6 +116,9 @@ def _rearrange_events_for_async_function_responses_in_history(
117116
function_call_id = function_response.id
118117
function_call_id_to_response_events_index[function_call_id] = i
119118

119+
if not function_call_id_to_response_events_index:
120+
return events
121+
120122
result_events: list[Event] = []
121123
for event in events:
122124
if event.get_function_responses():

src/google/adk/models/anthropic_llm.py

Lines changed: 7 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -553,11 +553,14 @@ async def generate_content_async(
553553
else NOT_GIVEN
554554
)
555555
thinking = _build_anthropic_thinking_param(llm_request.config)
556+
system = NOT_GIVEN
557+
if llm_request.config.system_instruction is not None:
558+
system = llm_request.config.system_instruction
556559

557560
if not stream:
558561
message = await self._anthropic_client.messages.create(
559562
model=model_to_use,
560-
system=llm_request.config.system_instruction,
563+
system=system,
561564
messages=messages,
562565
tools=tools,
563566
tool_choice=tool_choice,
@@ -567,14 +570,15 @@ async def generate_content_async(
567570
yield message_to_generate_content_response(message)
568571
else:
569572
async for response in self._generate_content_streaming(
570-
llm_request, messages, tools, tool_choice, thinking
573+
llm_request, messages, system, tools, tool_choice, thinking
571574
):
572575
yield response
573576

574577
async def _generate_content_streaming(
575578
self,
576579
llm_request: LlmRequest,
577580
messages: list[anthropic_types.MessageParam],
581+
system: Union[str, types.Content, NotGiven],
578582
tools: Union[Iterable[anthropic_types.ToolUnionParam], NotGiven],
579583
tool_choice: Union[anthropic_types.ToolChoiceParam, NotGiven],
580584
thinking: Union[
@@ -591,7 +595,7 @@ async def _generate_content_streaming(
591595
model_to_use = self._resolve_model_name(llm_request.model)
592596
raw_stream = await self._anthropic_client.messages.create(
593597
model=model_to_use,
594-
system=llm_request.config.system_instruction,
598+
system=system,
595599
messages=messages,
596600
tools=tools,
597601
tool_choice=tool_choice,

tests/unittests/a2a/executor/test_a2a_agent_executor.py

Lines changed: 76 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1072,3 +1072,79 @@ async def mock_run_async(**kwargs):
10721072
assert (
10731073
modified_a2a_event in enqueued_events
10741074
), "The modified event should have been enqueued"
1075+
1076+
@pytest.mark.asyncio
1077+
async def test_handle_request_preserves_metadata_in_final_events(
1078+
self,
1079+
) -> None:
1080+
"""Test that final events preserve invocation_id, author, and event_id in metadata."""
1081+
# Setup context with task_id
1082+
self.mock_context.task_id = "test-task-id"
1083+
self.mock_context.context_id = "test-context-id"
1084+
1085+
# Setup detailed mocks
1086+
self.mock_request_converter.return_value = AgentRunRequest(
1087+
user_id="test-user",
1088+
session_id="test-session",
1089+
new_message=Mock(spec=Content),
1090+
run_config=Mock(spec=RunConfig),
1091+
)
1092+
1093+
# Mock session service
1094+
mock_session = Mock()
1095+
mock_session.id = "test-session"
1096+
self.mock_runner.session_service.get_session = AsyncMock(
1097+
return_value=mock_session
1098+
)
1099+
1100+
# Mock invocation context
1101+
mock_invocation_context = Mock()
1102+
self.mock_runner._new_invocation_context.return_value = (
1103+
mock_invocation_context
1104+
)
1105+
1106+
# Mock ADK event with specific metadata to preserve
1107+
mock_adk_event = Mock(spec=Event)
1108+
mock_adk_event.invocation_id = "test-invocation-id"
1109+
mock_adk_event.author = "test-author"
1110+
mock_adk_event.id = "test-event-id"
1111+
1112+
# Configure run_async to yield our mock ADK event
1113+
async def mock_run_async(**kwargs):
1114+
async for item in self._create_async_generator([mock_adk_event]):
1115+
yield item
1116+
1117+
self.mock_runner.run_async = mock_run_async
1118+
self.mock_event_converter.return_value = [Mock()]
1119+
1120+
with patch(
1121+
"google.adk.a2a.executor.a2a_agent_executor.TaskResultAggregator"
1122+
) as mock_aggregator_class:
1123+
mock_aggregator = Mock()
1124+
mock_aggregator.task_state = TaskState.completed
1125+
mock_aggregator.task_status_message = Mock(spec=Message)
1126+
mock_aggregator_class.return_value = mock_aggregator
1127+
1128+
# Execute
1129+
await self.executor._handle_request(
1130+
self.mock_context, self.mock_event_queue
1131+
)
1132+
1133+
# Verify final status event was published and has correct metadata
1134+
final_events = [
1135+
call[0][0]
1136+
for call in self.mock_event_queue.enqueue_event.call_args_list
1137+
if hasattr(call[0][0], "final") and call[0][0].final == True
1138+
]
1139+
assert len(final_events) >= 1
1140+
final_event = final_events[-1]
1141+
1142+
assert final_event.metadata is not None
1143+
assert (
1144+
final_event.metadata.get("adk_invocation_id") == "test-invocation-id"
1145+
)
1146+
assert final_event.metadata.get("adk_author") == "test-author"
1147+
assert final_event.metadata.get("adk_event_id") == "test-event-id"
1148+
assert final_event.metadata.get("adk_app_name") == "test-app"
1149+
assert final_event.metadata.get("adk_user_id") == "test-user"
1150+
assert final_event.metadata.get("adk_session_id") == "test-session"

tests/unittests/flows/llm_flows/test_contents.py

Lines changed: 34 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1306,3 +1306,37 @@ def test_get_contents_live_history_rebuild():
13061306

13071307
assert result[1].role == "user"
13081308
assert "returned result" in result[1].parts[1].text
1309+
1310+
1311+
def test_rearrange_async_function_responses_early_returns_when_no_responses():
1312+
"""Rearrangement is a no-op when no event carries function_responses."""
1313+
events = [
1314+
Event(
1315+
invocation_id="inv1",
1316+
author="user",
1317+
content=types.UserContent("hi"),
1318+
),
1319+
Event(
1320+
invocation_id="inv2",
1321+
author="test_agent",
1322+
content=types.ModelContent("hello"),
1323+
),
1324+
Event(
1325+
invocation_id="inv3",
1326+
author="test_agent",
1327+
content=types.Content(
1328+
role="model",
1329+
parts=[
1330+
types.Part(
1331+
function_call=types.FunctionCall(
1332+
id="adk-1", name="tool", args={}
1333+
)
1334+
)
1335+
],
1336+
),
1337+
),
1338+
]
1339+
result = contents._rearrange_events_for_async_function_responses_in_history( # pylint: disable=protected-access
1340+
events
1341+
)
1342+
assert result is events

tests/unittests/models/test_anthropic_llm.py

Lines changed: 91 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@
2121
from unittest.mock import AsyncMock
2222
from unittest.mock import MagicMock
2323

24+
from anthropic import NOT_GIVEN
2425
from anthropic import types as anthropic_types
2526
from google.adk import version as adk_version
2627
from google.adk.models import anthropic_llm
@@ -2134,3 +2135,93 @@ async def test_generate_content_async_pairs_invalid_tool_ids(
21342135
]
21352136
assert len(set(use_ids)) == expected_unique
21362137
assert set(use_ids) == set(result_ids)
2138+
2139+
2140+
@pytest.mark.asyncio
2141+
async def test_non_streaming_no_system_instruction_passes_not_given():
2142+
"""system=NOT_GIVEN when LlmRequest has no system_instruction."""
2143+
llm = AnthropicLlm(model="claude-sonnet-4-20250514")
2144+
2145+
mock_message = anthropic_types.Message(
2146+
id="msg_test",
2147+
content=[
2148+
anthropic_types.TextBlock(text="ok", type="text", citations=None)
2149+
],
2150+
model="claude-sonnet-4-20250514",
2151+
role="assistant",
2152+
stop_reason="end_turn",
2153+
stop_sequence=None,
2154+
type="message",
2155+
usage=anthropic_types.Usage(
2156+
input_tokens=1,
2157+
output_tokens=1,
2158+
cache_creation_input_tokens=0,
2159+
cache_read_input_tokens=0,
2160+
server_tool_use=None,
2161+
service_tier=None,
2162+
),
2163+
)
2164+
2165+
mock_client = MagicMock()
2166+
mock_client.messages.create = AsyncMock(return_value=mock_message)
2167+
2168+
request = LlmRequest(
2169+
model="claude-sonnet-4-20250514",
2170+
contents=[Content(role="user", parts=[Part.from_text(text="Hi")])],
2171+
)
2172+
assert request.config.system_instruction is None
2173+
2174+
with mock.patch.object(llm, "_anthropic_client", mock_client):
2175+
_ = [r async for r in llm.generate_content_async(request, stream=False)]
2176+
2177+
mock_client.messages.create.assert_called_once()
2178+
_, kwargs = mock_client.messages.create.call_args
2179+
assert kwargs["system"] is NOT_GIVEN
2180+
2181+
2182+
@pytest.mark.asyncio
2183+
async def test_streaming_no_system_instruction_passes_not_given():
2184+
"""system=NOT_GIVEN on the streaming path when no system_instruction."""
2185+
llm = AnthropicLlm(model="claude-sonnet-4-20250514")
2186+
2187+
events = [
2188+
MagicMock(
2189+
type="message_start",
2190+
message=MagicMock(usage=MagicMock(input_tokens=1, output_tokens=0)),
2191+
),
2192+
MagicMock(
2193+
type="content_block_start",
2194+
index=0,
2195+
content_block=anthropic_types.TextBlock(text="", type="text"),
2196+
),
2197+
MagicMock(
2198+
type="content_block_delta",
2199+
index=0,
2200+
delta=anthropic_types.TextDelta(text="ok", type="text_delta"),
2201+
),
2202+
MagicMock(type="content_block_stop", index=0),
2203+
MagicMock(
2204+
type="message_delta",
2205+
delta=MagicMock(stop_reason="end_turn"),
2206+
usage=MagicMock(output_tokens=1),
2207+
),
2208+
MagicMock(type="message_stop"),
2209+
]
2210+
2211+
mock_client = MagicMock()
2212+
mock_client.messages.create = AsyncMock(
2213+
return_value=_make_mock_stream_events(events)
2214+
)
2215+
2216+
request = LlmRequest(
2217+
model="claude-sonnet-4-20250514",
2218+
contents=[Content(role="user", parts=[Part.from_text(text="Hi")])],
2219+
)
2220+
assert request.config.system_instruction is None
2221+
2222+
with mock.patch.object(llm, "_anthropic_client", mock_client):
2223+
_ = [r async for r in llm.generate_content_async(request, stream=True)]
2224+
2225+
mock_client.messages.create.assert_called_once()
2226+
_, kwargs = mock_client.messages.create.call_args
2227+
assert kwargs["system"] is NOT_GIVEN

0 commit comments

Comments
 (0)