From 81d0a3fd3cf1f809f962b224d62e71c098e9cdd8 Mon Sep 17 00:00:00 2001 From: mesutoezdil Date: Sat, 25 Jul 2026 12:22:26 +0200 Subject: [PATCH 1/2] fix(adk): keep full session state when num_recent_events is set get_session only fetched the last N events when num_recent_events was set, then built session.state by replaying just those N events. Any state_delta set by an older event outside that window never made it into state, silently. Now the full event history is always fetched to build state, and num_recent_events only trims the events list returned on the session, after state is already correct. Signed-off-by: mesutoezdil --- .../src/kagent/adk/_session_service.py | 29 ++++++------ .../tests/unittests/test_session_service.py | 46 +++++++++++++++++-- 2 files changed, 57 insertions(+), 18 deletions(-) diff --git a/python/packages/kagent-adk/src/kagent/adk/_session_service.py b/python/packages/kagent-adk/src/kagent/adk/_session_service.py index fd6ec74ea..820fedb9a 100644 --- a/python/packages/kagent-adk/src/kagent/adk/_session_service.py +++ b/python/packages/kagent-adk/src/kagent/adk/_session_service.py @@ -47,7 +47,7 @@ async def create_session( # Make API call to create session # Pass user_id as a query param so the controller's auth middleware - # (UnsecureAuthenticator) reads it consistently — matching the user_id + # (UnsecureAuthenticator) reads it consistently, matching the user_id # used by get_session, list_sessions, delete_session, and append_event. # Without this, unsecure-mode requests fall back to "admin@kagent.dev" # while all lookups use the A2A-derived user_id, causing SessionNotFoundError. @@ -77,20 +77,16 @@ async def get_session( config: Optional[GetSessionConfig] = None, ) -> Optional[Session]: try: - # ADK requires events to be chronological (especially for calculating deltas) - url = f"/api/sessions/{session_id}?user_id={user_id}&order=asc" - if config: - if config.after_timestamp: - # TODO: implement - # url += f"&after={config.after_timestamp}" - pass - if config.num_recent_events: - url += f"&limit={config.num_recent_events}" - else: - url += "&limit=-1" - else: - # return all - url += "&limit=-1" + # ADK requires events to be chronological (especially for calculating deltas). + # Always fetch the full history: state is built by replaying every event's + # state_delta below, so limiting the fetch here would silently drop state set + # by events outside the window. num_recent_events is applied after, by + # trimming session.events once state is already correct. + url = f"/api/sessions/{session_id}?user_id={user_id}&order=asc&limit=-1" + if config and config.after_timestamp: + # TODO: implement + # url += f"&after={config.after_timestamp}" + pass # Make API call to get session response: httpx.Response = await self.client.get(url) @@ -124,6 +120,9 @@ async def get_session( for event in events: await super().append_event(session, event) + if config and config.num_recent_events: + session.events = session.events[-config.num_recent_events :] + return session except httpx.HTTPStatusError as e: if e.response.status_code == 404: diff --git a/python/packages/kagent-adk/tests/unittests/test_session_service.py b/python/packages/kagent-adk/tests/unittests/test_session_service.py index fb9fce139..d5fc883c0 100644 --- a/python/packages/kagent-adk/tests/unittests/test_session_service.py +++ b/python/packages/kagent-adk/tests/unittests/test_session_service.py @@ -5,6 +5,7 @@ import httpx import pytest from google.adk.events.event import Event, EventActions +from google.adk.sessions.base_session_service import GetSessionConfig from kagent.adk._session_service import KAgentSessionService @@ -74,7 +75,7 @@ async def test_create_session_passes_user_id_as_query_param(): the controller's UnsecureAuthenticator resolves identity from the query param (or X-User-Id header), not the JSON body. Without the query param the controller falls back to "admin@kagent.dev" for the session create, while - every subsequent GET uses the A2A-derived user_id — guaranteeing a 404. + every subsequent GET uses the A2A-derived user_id, guaranteeing a 404. Fixes: https://github.com/kagent-dev/kagent/issues/1882 """ mock_response = MagicMock(spec=httpx.Response) @@ -138,7 +139,7 @@ async def test_get_session_events_not_duplicated(make_event, session_response, s assert session is not None assert len(session.events) == len(events), ( - f"Expected {len(events)} events but got {len(session.events)} — possible event duplication in get_session" + f"Expected {len(events)} events but got {len(session.events)}, possible event duplication in get_session" ) @@ -178,11 +179,50 @@ async def test_get_session_state_delta_applied_once(make_event, session_response # so for an idempotent string the bug was silent; here we use a distinct value # and just verify the key is present with the correct value.) assert session.state.get("counter") == 7, ( - f"Expected state['counter'] == 7, got {session.state.get('counter')} — " + f"Expected state['counter'] == 7, got {session.state.get('counter')}, " "state_delta may have been applied more than once" ) +@pytest.mark.asyncio +async def test_get_session_state_kept_outside_recent_events_window(make_event, session_response): + """A state delta from an event outside the num_recent_events window must + still land in session.state, only session.events is trimmed to the window. + + The mock server here behaves like the real one: a &limit=N query param + returns only the last N events, while &limit=-1 returns everything. This + is what makes the test fail against the old code, which asked the server + for only num_recent_events and so never saw the older state_delta at all. + """ + all_events = [ + make_event("assistant", state_delta={"old_key": "old_value"}), + make_event("user"), + make_event("assistant"), + ] + + def get_side_effect(url: str): + limited = "limit=-1" not in url + events = all_events[-2:] if limited else all_events + mock_response = MagicMock(spec=httpx.Response) + mock_response.status_code = 200 + mock_response.json.return_value = session_response(events) + mock_response.raise_for_status = MagicMock() + return mock_response + + client = MagicMock(spec=httpx.AsyncClient) + client.get = AsyncMock(side_effect=get_side_effect) + + session = await KAgentSessionService(client).get_session( + app_name="app", user_id="u1", session_id="s1", config=GetSessionConfig(num_recent_events=2) + ) + + assert session is not None + assert len(session.events) == 2 + assert session.state.get("old_key") == "old_value", ( + "state from the first event must still apply even though only the last 2 events are kept in session.events" + ) + + @pytest.mark.asyncio async def test_get_session_multiple_state_deltas_applied_once(make_event, session_response, service): """Multiple events each contributing a state key are each applied once.""" From 9c018d94ad23e3038a139831b4a0fc75ee493f89 Mon Sep 17 00:00:00 2001 From: mesutoezdil Date: Tue, 28 Jul 2026 23:29:56 +0200 Subject: [PATCH 2/2] fix(adk): rebuild state before trimming events for num_recent_events=0 num_recent_events=0 emptied the event list before the replay loop, so the session came back with no state at all. Trim after the replay instead, like the other window sizes, and guard the [-0:] case which would keep everything. Signed-off-by: mesutoezdil --- .../kagent-adk/src/kagent/adk/_session_service.py | 10 ++++++---- .../kagent-adk/tests/unittests/test_session_service.py | 8 ++++++-- 2 files changed, 12 insertions(+), 6 deletions(-) diff --git a/python/packages/kagent-adk/src/kagent/adk/_session_service.py b/python/packages/kagent-adk/src/kagent/adk/_session_service.py index 659475c01..e3d08e4f3 100644 --- a/python/packages/kagent-adk/src/kagent/adk/_session_service.py +++ b/python/packages/kagent-adk/src/kagent/adk/_session_service.py @@ -102,8 +102,6 @@ async def get_session( session_data = data["data"]["session"] events_data = data["data"]["events"] - if config and config.num_recent_events == 0: - events_data = [] events: list[Event] = [] for event_data in events_data: @@ -121,8 +119,12 @@ async def get_session( for event in events: await super().append_event(session, event) - if config and config.num_recent_events: - session.events = session.events[-config.num_recent_events :] + if config and config.num_recent_events is not None: + # Trim only after every event has been replayed, so state is complete. + # num_recent_events == 0 means "no events", not "all events" ([-0:] would + # keep everything). + num_recent_events = config.num_recent_events + session.events = session.events[-num_recent_events:] if num_recent_events else [] return session except httpx.HTTPStatusError as e: diff --git a/python/packages/kagent-adk/tests/unittests/test_session_service.py b/python/packages/kagent-adk/tests/unittests/test_session_service.py index ff99f6684..6655a6dd5 100644 --- a/python/packages/kagent-adk/tests/unittests/test_session_service.py +++ b/python/packages/kagent-adk/tests/unittests/test_session_service.py @@ -165,8 +165,11 @@ async def test_get_session_passes_epoch_timestamp_to_api(mock_client, session_re @pytest.mark.asyncio async def test_get_session_with_zero_recent_events_returns_no_events(make_event, session_response, mock_client): - """ADK defines a zero recent-event limit as returning session metadata without history.""" - client = mock_client(session_response([make_event("user")])) + """ADK defines a zero recent-event limit as returning session metadata without history. + + State still has to be complete: every event is replayed, only the events list is emptied. + """ + client = mock_client(session_response([make_event("user", state_delta={"key": "value"})])) svc = KAgentSessionService(client) session = await svc.get_session( @@ -178,6 +181,7 @@ async def test_get_session_with_zero_recent_events_returns_no_events(make_event, assert session is not None assert session.events == [] + assert session.state.get("key") == "value", "state must survive even when no events are returned" client.get.assert_awaited_once_with( "/api/sessions/s1", params={"user_id": "u1", "order": "asc", "limit": -1},