diff --git a/constraints-3.10.txt b/constraints-3.10.txt index 59bdb41aee7..034b27819c1 100644 --- a/constraints-3.10.txt +++ b/constraints-3.10.txt @@ -1,5 +1,5 @@ # This file was autogenerated by uv via the following command: -# uv pip compile pyproject.toml --all-extras --python-version 3.10 --no-emit-package google-adk --exclude-newer 2026-08-20 --index-url https://pypi.org/simple -o constraints-3.10.txt +# uv pip compile pyproject.toml --all-extras --python-version 3.10 --no-emit-package google-adk --exclude-newer 2026-08-25 --index-url https://pypi.org/simple -o constraints-3.10.txt a2a-sdk==1.1.2 # via # -c constraints-3.10.txt.stable.tmp @@ -147,6 +147,15 @@ black==25.12.0 # via # -c constraints-3.10.txt.stable.tmp # pyink +boto3==1.43.80 + # via + # -c constraints-3.10.txt.stable.tmp + # google-adk (pyproject.toml) +botocore==1.43.80 + # via + # -c constraints-3.10.txt.stable.tmp + # boto3 + # s3transfer bracex==3.0.1 # via # -c constraints-3.10.txt.stable.tmp @@ -792,6 +801,11 @@ jiter==0.16.0 # -c constraints-3.10.txt.stable.tmp # anthropic # openai +jmespath==1.1.0 + # via + # -c constraints-3.10.txt.stable.tmp + # boto3 + # botocore joblib==1.5.3 # via # -c constraints-3.10.txt.stable.tmp @@ -1430,6 +1444,7 @@ python-dateutil==2.9.0.post0 # via # -c constraints-3.10.txt.stable.tmp # google-adk (pyproject.toml) + # botocore # daytona-analytics-api-client # daytona-analytics-api-client-async # daytona-api-client @@ -1565,6 +1580,10 @@ ruff==0.15.17 # via # -c constraints-3.10.txt.stable.tmp # google-adk (pyproject.toml) +s3transfer==0.19.2 + # via + # -c constraints-3.10.txt.stable.tmp + # boto3 scikit-learn==1.5.2 # via # -c constraints-3.10.txt.stable.tmp @@ -1885,6 +1904,7 @@ uritemplate==4.2.0 urllib3==2.7.0 # via # -c constraints-3.10.txt.stable.tmp + # botocore # daytona # daytona-analytics-api-client # daytona-api-client diff --git a/constraints-3.11.txt b/constraints-3.11.txt index 2f6ab89f92b..03e7c4c0302 100644 --- a/constraints-3.11.txt +++ b/constraints-3.11.txt @@ -1,5 +1,5 @@ # This file was autogenerated by uv via the following command: -# uv pip compile pyproject.toml --all-extras --python-version 3.11 --no-emit-package google-adk --exclude-newer 2026-08-20 --index-url https://pypi.org/simple -o constraints-3.11.txt +# uv pip compile pyproject.toml --all-extras --python-version 3.11 --no-emit-package google-adk --exclude-newer 2026-08-25 --index-url https://pypi.org/simple -o constraints-3.11.txt a2a-sdk==1.1.1 # via # -c constraints-3.11.txt.stable.tmp @@ -161,6 +161,15 @@ black==25.12.0 # via # -c constraints-3.11.txt.stable.tmp # pyink +boto3==1.43.80 + # via + # -c constraints-3.11.txt.stable.tmp + # google-adk (pyproject.toml) +botocore==1.43.80 + # via + # -c constraints-3.11.txt.stable.tmp + # boto3 + # s3transfer bracex==3.0.1 # via # -c constraints-3.11.txt.stable.tmp @@ -876,6 +885,8 @@ jiter==0.14.0 jmespath==1.1.0 # via # -c constraints-3.11.txt.stable.tmp + # boto3 + # botocore # cel-python joblib==1.5.3 # via @@ -1652,6 +1663,7 @@ python-dateutil==2.9.0.post0 # via # -c constraints-3.11.txt.stable.tmp # google-adk (pyproject.toml) + # botocore # daytona-analytics-api-client # daytona-analytics-api-client-async # daytona-api-client @@ -1823,6 +1835,10 @@ ruff==0.15.17 # via # -c constraints-3.11.txt.stable.tmp # google-adk (pyproject.toml) +s3transfer==0.19.2 + # via + # -c constraints-3.11.txt.stable.tmp + # boto3 scikit-learn==1.9.0 # via # -c constraints-3.11.txt.stable.tmp @@ -2153,6 +2169,7 @@ uritemplate==4.2.0 urllib3==2.7.0 # via # -c constraints-3.11.txt.stable.tmp + # botocore # daytona # daytona-analytics-api-client # daytona-api-client diff --git a/constraints-3.12.txt b/constraints-3.12.txt index 0d49c542718..52c071883ad 100644 --- a/constraints-3.12.txt +++ b/constraints-3.12.txt @@ -1,5 +1,5 @@ # This file was autogenerated by uv via the following command: -# uv pip compile pyproject.toml --all-extras --python-version 3.12 --no-emit-package google-adk --exclude-newer 2026-08-20 --index-url https://pypi.org/simple -o constraints-3.12.txt +# uv pip compile pyproject.toml --all-extras --python-version 3.12 --no-emit-package google-adk --exclude-newer 2026-08-25 --index-url https://pypi.org/simple -o constraints-3.12.txt a2a-sdk==1.1.1 # via # -c constraints-3.12.txt.stable.tmp @@ -137,6 +137,15 @@ black==25.12.0 # via # -c constraints-3.12.txt.stable.tmp # pyink +boto3==1.43.80 + # via + # -c constraints-3.12.txt.stable.tmp + # google-adk (pyproject.toml) +botocore==1.43.80 + # via + # -c constraints-3.12.txt.stable.tmp + # boto3 + # s3transfer bracex==3.0.1 # via # -c constraints-3.12.txt.stable.tmp @@ -776,6 +785,11 @@ jiter==0.16.0 # -c constraints-3.12.txt.stable.tmp # anthropic # openai +jmespath==1.1.0 + # via + # -c constraints-3.12.txt.stable.tmp + # boto3 + # botocore joblib==1.5.3 # via # -c constraints-3.12.txt.stable.tmp @@ -1418,6 +1432,7 @@ python-dateutil==2.9.0.post0 # via # -c constraints-3.12.txt.stable.tmp # google-adk (pyproject.toml) + # botocore # daytona-analytics-api-client # daytona-analytics-api-client-async # daytona-api-client @@ -1561,6 +1576,10 @@ ruff==0.15.17 # via # -c constraints-3.12.txt.stable.tmp # google-adk (pyproject.toml) +s3transfer==0.19.2 + # via + # -c constraints-3.12.txt.stable.tmp + # boto3 scikit-learn==1.9.0 # via # -c constraints-3.12.txt.stable.tmp @@ -1852,6 +1871,7 @@ uritemplate==4.2.0 urllib3==2.7.0 # via # -c constraints-3.12.txt.stable.tmp + # botocore # daytona # daytona-analytics-api-client # daytona-api-client diff --git a/constraints-3.13.txt b/constraints-3.13.txt index 48563bbc6ee..36fde57efb4 100644 --- a/constraints-3.13.txt +++ b/constraints-3.13.txt @@ -1,5 +1,5 @@ # This file was autogenerated by uv via the following command: -# uv pip compile pyproject.toml --all-extras --python-version 3.13 --no-emit-package google-adk --exclude-newer 2026-08-20 --index-url https://pypi.org/simple -o constraints-3.13.txt +# uv pip compile pyproject.toml --all-extras --python-version 3.13 --no-emit-package google-adk --exclude-newer 2026-08-25 --index-url https://pypi.org/simple -o constraints-3.13.txt a2a-sdk==1.1.1 # via # -c constraints-3.13.txt.stable.tmp @@ -133,6 +133,15 @@ black==25.12.0 # via # -c constraints-3.13.txt.stable.tmp # pyink +boto3==1.43.80 + # via + # -c constraints-3.13.txt.stable.tmp + # google-adk (pyproject.toml) +botocore==1.43.80 + # via + # -c constraints-3.13.txt.stable.tmp + # boto3 + # s3transfer bracex==3.0.1 # via # -c constraints-3.13.txt.stable.tmp @@ -768,6 +777,11 @@ jiter==0.16.0 # -c constraints-3.13.txt.stable.tmp # anthropic # openai +jmespath==1.1.0 + # via + # -c constraints-3.13.txt.stable.tmp + # boto3 + # botocore joblib==1.5.3 # via # -c constraints-3.13.txt.stable.tmp @@ -1410,6 +1424,7 @@ python-dateutil==2.9.0.post0 # via # -c constraints-3.13.txt.stable.tmp # google-adk (pyproject.toml) + # botocore # daytona-analytics-api-client # daytona-analytics-api-client-async # daytona-api-client @@ -1553,6 +1568,10 @@ ruff==0.15.17 # via # -c constraints-3.13.txt.stable.tmp # google-adk (pyproject.toml) +s3transfer==0.19.2 + # via + # -c constraints-3.13.txt.stable.tmp + # boto3 scikit-learn==1.9.0 # via # -c constraints-3.13.txt.stable.tmp @@ -1833,6 +1852,7 @@ uritemplate==4.2.0 urllib3==2.7.0 # via # -c constraints-3.13.txt.stable.tmp + # botocore # daytona # daytona-analytics-api-client # daytona-api-client diff --git a/constraints-3.14.txt b/constraints-3.14.txt index 97f66e25a94..778789741ac 100644 --- a/constraints-3.14.txt +++ b/constraints-3.14.txt @@ -1,5 +1,5 @@ # This file was autogenerated by uv via the following command: -# uv pip compile pyproject.toml --all-extras --python-version 3.14 --no-emit-package google-adk --exclude-newer 2026-08-20 --index-url https://pypi.org/simple -o constraints-3.14.txt +# uv pip compile pyproject.toml --all-extras --python-version 3.14 --no-emit-package google-adk --exclude-newer 2026-08-25 --index-url https://pypi.org/simple -o constraints-3.14.txt a2a-sdk==1.1.1 # via # -c constraints-3.14.txt.stable.tmp @@ -133,6 +133,15 @@ black==25.12.0 # via # -c constraints-3.14.txt.stable.tmp # pyink +boto3==1.43.80 + # via + # -c constraints-3.14.txt.stable.tmp + # google-adk (pyproject.toml) +botocore==1.43.80 + # via + # -c constraints-3.14.txt.stable.tmp + # boto3 + # s3transfer bracex==3.0.1 # via # -c constraints-3.14.txt.stable.tmp @@ -768,6 +777,11 @@ jiter==0.16.0 # -c constraints-3.14.txt.stable.tmp # anthropic # openai +jmespath==1.1.0 + # via + # -c constraints-3.14.txt.stable.tmp + # boto3 + # botocore joblib==1.5.3 # via # -c constraints-3.14.txt.stable.tmp @@ -1410,6 +1424,7 @@ python-dateutil==2.9.0.post0 # via # -c constraints-3.14.txt.stable.tmp # google-adk (pyproject.toml) + # botocore # daytona-analytics-api-client # daytona-analytics-api-client-async # daytona-api-client @@ -1553,6 +1568,10 @@ ruff==0.15.17 # via # -c constraints-3.14.txt.stable.tmp # google-adk (pyproject.toml) +s3transfer==0.19.2 + # via + # -c constraints-3.14.txt.stable.tmp + # boto3 scikit-learn==1.9.0 # via # -c constraints-3.14.txt.stable.tmp @@ -1833,6 +1852,7 @@ uritemplate==4.2.0 urllib3==2.7.0 # via # -c constraints-3.14.txt.stable.tmp + # botocore # daytona # daytona-analytics-api-client # daytona-api-client diff --git a/docs/guides/README.md b/docs/guides/README.md index 249165eec78..5a7da42f879 100644 --- a/docs/guides/README.md +++ b/docs/guides/README.md @@ -31,6 +31,8 @@ This directory contains specific developer guides for the ADK Python implementat * [Live model callbacks](flows/llm_flows/base_llm_flow/live_model_callbacks.md) - Inspecting or blocking content on a live bidirectional session. ### Integrations +* [AgentCoreSessionService](integrations/agentcore/agentcore_session_service/index.md) - Storing ADK sessions in Amazon Bedrock AgentCore Memory short-term memory. +* [AgentCoreSessionServiceConfig](integrations/agentcore/config/index.md) - Memory id and AWS region for AgentCoreSessionService. * [Model Armor](integrations/model_armor/index.md) - Screening user input and model output with Google Cloud Model Armor. ### Labs diff --git a/docs/guides/integrations/agentcore/agentcore_session_service/index.md b/docs/guides/integrations/agentcore/agentcore_session_service/index.md new file mode 100644 index 00000000000..984732349bd --- /dev/null +++ b/docs/guides/integrations/agentcore/agentcore_session_service/index.md @@ -0,0 +1,120 @@ +# AgentCoreSessionService + +`AgentCoreSessionService` is a `BaseSessionService` that stores ADK sessions in Amazon Bedrock AgentCore Memory short-term memory. Pass it to `Runner` as `session_service` when the conversation history should live in AgentCore instead of in memory, Redis, or Firestore. + +## Introduction + +AgentCore Memory has no session resource. An actor and a session id identify a stream of events, and each event is one turn payload. ADK sessions are finer-grained: they carry a list of `Event` objects, session state, and an `(app_name, user_id, session_id)` key. + +`AgentCoreSessionService` is the compatibility layer. `Runner` depends on `BaseSessionService`, so any agent that already uses `Runner(..., session_service=...)` can swap this implementation in. The service writes each non-partial ADK event as one AgentCore event: conversational text so AgentCore can extract long-term memories from the turn, and a blob of the full ADK event JSON so tool calls, state deltas, and metadata come back on `get_session`. + +Install the extra, create an AgentCore Memory resource, and point the service at its id. AWS credentials come from the boto3 default chain. + +## Get started + +```python +from google.adk.agents import Agent +from google.adk.integrations.agentcore import AgentCoreSessionService +from google.adk.runners import Runner + +session_service = AgentCoreSessionService( + memory_id="my-memory-abc123", + region_name="us-east-1", +) + +agent = Agent( + name="assistant", + instruction="You are a helpful AI assistant.", +) + +runner = Runner( + app_name="my_app", + agent=agent, + session_service=session_service, +) +``` + +```shell +pip install 'google-adk[agentcore]' +``` + +`create_session` / `get_session` / `append_event` work the same as on `InMemorySessionService`. `Runner.run_async` loads the session, appends user and model events, and the service flushes each complete event to AgentCore. + +## How it works + +ADK identifies a session by `(app_name, user_id, session_id)`. AgentCore identifies a stream by `(memoryId, actorId, sessionId)`. The service maps: + +- `actorId` to `{app_name}:{user_id}`, so two ADK apps sharing one Memory resource do not mix users +- `sessionId` to the ADK session id +- each non-partial `Event` to one `create_event` call + +When the event has text parts, the AgentCore payload starts with a `conversational` item whose role is `USER`, `ASSISTANT`, or `TOOL`. A `blob` of `Event.model_dump_json()` always follows, and `get_session` rebuilds history from those blobs. Streaming fragments with `partial=True` are not written, matching other persistent backends. + +AgentCore has no create-session or delete-session API. `create_session` writes a bootstrap event with `extractionMode=SKIP` so the session appears in `list_sessions` before the first user turn. `delete_session` deletes every event in that stream. + +`list_sessions` with a `user_id` lists that actor. Omitting `user_id` lists every actor whose id starts with `{app_name}:`. Returned sessions have empty `events` lists, as `BaseSessionService` requires. + +## Configuration options + +### Constructor + +| Option | Type | Default | Description | +| :--- | :--- | :--- | :--- | +| `memory_id` | `str \| None` | `None` | AgentCore Memory resource id. Required when `config` is omitted. | +| `config` | `AgentCoreSessionServiceConfig \| None` | `None` | Full config object. Overrides `memory_id` and `region_name` when set. | +| `client` | boto3 client or test double | `None` | Pre-built `bedrock-agentcore` client. Built lazily from boto3 when omitted. | +| `region_name` | `str \| None` | `None` | AWS region. Ignored when `config` is set. Falls back to the boto3 default chain. | + +- **`memory_id`** is the Memory resource you created in AgentCore. Every event is written under this id. +- **`config`** is useful when you already constructed `AgentCoreSessionServiceConfig`. See that type for field-level detail. +- **`client`** is how tests inject a fake, and how an application can pass a client with custom retries or endpoints. +- **`region_name`** is needed when the default AWS region is not the region where the Memory resource lives. + +### `AgentCoreSessionServiceConfig` fields + +| Option | Type | Default | Description | +| :--- | :--- | :--- | :--- | +| `memory_id` | `str` | required | AgentCore Memory resource id. | +| `region_name` | `str \| None` | `None` | AWS region for the `bedrock-agentcore` client. | + +## Advanced applications + +### Passing a boto3 client + +```python +import boto3 +from google.adk.integrations.agentcore import AgentCoreSessionService + +client = boto3.client("bedrock-agentcore", region_name="eu-west-2") +session_service = AgentCoreSessionService( + memory_id="my-memory-abc123", + client=client, +) +``` + +Use this when credential loading, retries, or the endpoint should not follow process-wide boto3 defaults. + +### Filtering history on get + +```python +from google.adk.sessions.base_session_service import GetSessionConfig + +session = await session_service.get_session( + app_name="my_app", + user_id="user-42", + session_id="session-a", + config=GetSessionConfig(num_recent_events=20), +) +``` + +`after_timestamp` keeps events whose ADK `timestamp` is at least that unix time. `num_recent_events=0` returns the session with an empty event list. + +## Limitations + +- **App- and user-scoped state is not shared across sessions.** Redis and Firestore keep `app:` and `user:` keys in separate documents. This backend replays `state_delta` onto the session that stored the events. A new session for the same user does not inherit `user:` keys from older sessions. + +- **Empty AgentCore sessions expire.** AgentCore deletes empty sessions after about a day. `create_session` writes a bootstrap event so the session is not empty. If every event is later deleted, the session disappears. + +- **Session ids must be acceptable to AgentCore.** The service does not rewrite ids. Prefer the generated UUID, or a value in the character set AgentCore accepts. + +- **This is short-term memory only.** Long-term extraction is whatever the Memory resource is configured to do with conversational payloads. The ADK `BaseMemoryService` search API is a different interface. diff --git a/docs/guides/integrations/agentcore/config/index.md b/docs/guides/integrations/agentcore/config/index.md new file mode 100644 index 00000000000..806c6336e87 --- /dev/null +++ b/docs/guides/integrations/agentcore/config/index.md @@ -0,0 +1,40 @@ +# AgentCoreSessionServiceConfig + +`AgentCoreSessionServiceConfig` holds the AgentCore Memory resource id and optional AWS region for `AgentCoreSessionService`. Construct it when you want to pass configuration as an object rather than as `memory_id=` / `region_name=` keyword arguments. + +## Introduction + +`AgentCoreSessionService` needs to know which Memory resource to write to, and which AWS region hosts that resource. Those two values are the only fields on this config type. `AgentCoreSessionService(config=...)` takes precedence over the convenience `memory_id` and `region_name` constructor arguments. + +Callers who only have a memory id can skip this type and pass `memory_id=` directly. The config object is the better fit when configuration is loaded from a file or built in one place and handed to the service in another. + +## Get started + +```python +from google.adk.integrations.agentcore import AgentCoreSessionService +from google.adk.integrations.agentcore import AgentCoreSessionServiceConfig + +config = AgentCoreSessionServiceConfig( + memory_id="my-memory-abc123", + region_name="us-east-1", +) +session_service = AgentCoreSessionService(config=config) +``` + +## How it works + +The service reads `config.memory_id` on every AgentCore API call as `memoryId`. When `client` is omitted, the service builds a `boto3.client("bedrock-agentcore", ...)` and passes `config.region_name` if it is set. When `region_name` is unset, boto3 uses its default credential and region chain. + +## Configuration options + +| Option | Type | Default | Description | +| :--- | :--- | :--- | :--- | +| `memory_id` | `str` | required | AgentCore Memory resource id sent as `memoryId`. | +| `region_name` | `str \| None` | `None` | AWS region for the lazily built boto3 client. | + +- **`memory_id`** is the id of the Memory resource, not an ARN. Every `create_event`, `list_events`, `list_sessions`, `list_actors`, and `delete_event` call includes it. +- **`region_name`** is ignored when you pass a pre-built `client` to `AgentCoreSessionService`, because that client already has a region. + +## Related samples + +See [AgentCoreSessionService](../agentcore_session_service/index.md) for Runner wiring, event mapping, and limitations. diff --git a/pyproject.toml b/pyproject.toml index 0423900dc04..3d51b961cdc 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -64,6 +64,9 @@ optional-dependencies.agent-identity = [ "google-cloud-agentidentitycredentials>=0.1,<0.2", "google-cloud-iamconnectorcredentials>=0.1,<0.2", ] +optional-dependencies.agentcore = [ + "boto3>=1.43.77", # bedrock-agentcore Memory STM APIs (create_event / list_events). +] # Union of every extra that unlocks a runtime feature. Excludes benchmark, # community, dev, docs and test, which exist to build, test or document ADK # itself. @@ -72,6 +75,7 @@ optional-dependencies.all = [ "anthropic>=0.78", "anyio>=4.9,<5", "beautifulsoup4>=3.2.2", + "boto3>=1.43.77", "crewai[tools]; python_version>='3.11' and python_version<'3.12'", # chromadb/pypika fail on 3.12+ "daytona>=0.191", "docker>=7", diff --git a/src/google/adk/integrations/agentcore/README.md b/src/google/adk/integrations/agentcore/README.md new file mode 100644 index 00000000000..33e62685dfc --- /dev/null +++ b/src/google/adk/integrations/agentcore/README.md @@ -0,0 +1,89 @@ +# Copyright 2026 Google LLC +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +# AgentCore Memory Integration for ADK + +This integration stores ADK sessions in +[Amazon Bedrock AgentCore Memory](https://docs.aws.amazon.com/bedrock-agentcore/latest/devguide/memory.html) +short-term memory. + +AgentCore has no session object. Each ADK session is an AgentCore +`(actorId, sessionId)` pair, and each ADK `Event` is one AgentCore event. + +The mapping follows the approach in +[divakaivan/agentcore_session_service](https://github.com/divakaivan/agentcore_session_service) +(#6920): a conversational payload so AgentCore can extract memories from the +turn text, plus a blob of the full ADK event JSON so tool calls, state deltas, +and metadata round-trip. + +## Installation + +```bash +pip install "google-adk[agentcore]" +``` + +AWS credentials must be available to boto3 (environment, shared config, or +instance role). You also need an AgentCore Memory resource id. + +## Quick Start + +```python +from google.adk.agents import Agent +from google.adk.integrations.agentcore import AgentCoreSessionService +from google.adk.runners import Runner + +session_service = AgentCoreSessionService( + memory_id="my-memory-abc123", + region_name="us-east-1", +) + +agent = Agent( + name="assistant", + instruction="You are a helpful AI assistant.", +) + +runner = Runner( + app_name="my_app", + agent=agent, + session_service=session_service, +) +``` + +## Mapping + +| ADK | AgentCore | +| :--- | :--- | +| `(app_name, user_id)` | `actorId` = `{app_name}:{user_id}` | +| `session.id` | `sessionId` | +| `Event` (non-partial) | `create_event` payload: optional `conversational` text (`USER` / `ASSISTANT` / `TOOL`) and a `blob` of `Event.model_dump_json()` | + +`create_session` writes a bootstrap blob with `extractionMode=SKIP` so the +session exists in `list_sessions` / `get_session` before the first user turn. +`delete_session` deletes every event in that session (AgentCore has no +delete-session API). + +App- and user-scoped state is stored on the session via event `state_delta` +replay. It is not shared across sessions the way Redis/Firestore backends +share `app:` / `user:` keys. + +## Configuration + +`AgentCoreSessionServiceConfig`: + +| Field | Type | Default | Description | +| :--- | :--- | :--- | :--- | +| `memory_id` | `str` | required | AgentCore Memory resource id. | +| `region_name` | `Optional[str]` | `None` | AWS region. Falls back to the boto3 default chain. | + +You can pass a pre-built `boto3.client("bedrock-agentcore", ...)` as `client`. diff --git a/src/google/adk/integrations/agentcore/__init__.py b/src/google/adk/integrations/agentcore/__init__.py new file mode 100644 index 00000000000..2c2c44e76de --- /dev/null +++ b/src/google/adk/integrations/agentcore/__init__.py @@ -0,0 +1,25 @@ +# Copyright 2026 Google LLC +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +"""AWS Bedrock AgentCore Memory integrations for ADK.""" + +from __future__ import annotations + +from ._agentcore_session_service import AgentCoreSessionService +from ._config import AgentCoreSessionServiceConfig + +__all__ = [ + 'AgentCoreSessionService', + 'AgentCoreSessionServiceConfig', +] diff --git a/src/google/adk/integrations/agentcore/_agentcore_session_service.py b/src/google/adk/integrations/agentcore/_agentcore_session_service.py new file mode 100644 index 00000000000..3c9e92abc36 --- /dev/null +++ b/src/google/adk/integrations/agentcore/_agentcore_session_service.py @@ -0,0 +1,433 @@ +# Copyright 2026 Google LLC +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +"""AgentCore Memory-backed session service implementation for ADK.""" + +from __future__ import annotations + +import asyncio +from datetime import datetime +from datetime import timezone +import json +import logging +from typing import Any +from typing import Optional + +from ...errors.already_exists_error import AlreadyExistsError +from ...events.event import Event +from ...platform import uuid as platform_uuid +from ...sessions.base_session_service import BaseSessionService +from ...sessions.base_session_service import GetSessionConfig +from ...sessions.base_session_service import ListSessionsResponse +from ...sessions.session import Session +from ...sessions.state import State +from ...utils._dependency import missing_extra +from ._config import AgentCoreSessionServiceConfig + +logger = logging.getLogger('google_adk.' + __name__) + +# Stored as a blob so get_session can tell "this session exists" from +# "this session was never created". AgentCore has no CreateSession API. +_BOOTSTRAP_KIND = 'adk.session.bootstrap' +_ACTOR_SEPARATOR = ':' +_PAGE_SIZE = 100 + + +def _unix_timestamp(value: object) -> float: + """Converts an AgentCore eventTimestamp to a unix timestamp.""" + if isinstance(value, datetime): + if value.tzinfo is None: + value = value.replace(tzinfo=timezone.utc) + return value.timestamp() + if isinstance(value, (int, float)): + return float(value) + return 0.0 + + +class AgentCoreSessionService(BaseSessionService): + """Session service backed by AWS Bedrock AgentCore Memory short-term memory. + + AgentCore has no session resource. ADK sessions are mapped as: + + - ``actorId``: ``{app_name}:{user_id}`` + - ``sessionId``: the ADK session id + - each ADK ``Event``: one AgentCore event whose payload is the conversational + text (when present) plus a blob of the full ADK event JSON, so history can + be reconstructed without losing tool calls, state deltas, or metadata. + + Partial (streaming) events are not written. ``delete_session`` deletes every + AgentCore event for that session. + """ + + def __init__( + self, + *, + memory_id: Optional[str] = None, + config: Optional[AgentCoreSessionServiceConfig] = None, + client: Optional[Any] = None, + region_name: Optional[str] = None, + ): + """Initializes the AgentCoreSessionService. + + Args: + memory_id: AgentCore Memory resource id. Ignored when ``config`` is set. + config: Optional full config. If omitted, ``memory_id`` is required. + client: Optional pre-configured ``bedrock-agentcore`` boto3 client (or a + test double). When omitted, a client is created lazily. + region_name: AWS region. Ignored when ``config`` is set. + """ + if config is None: + if not memory_id: + raise ValueError('memory_id is required when config is not provided.') + config = AgentCoreSessionServiceConfig( + memory_id=memory_id, region_name=region_name + ) + self.config = config + self._client = client + + def _get_client(self) -> Any: + """Lazily initializes and returns the bedrock-agentcore client.""" + if self._client is not None: + return self._client + try: + import boto3 # type: ignore + except ImportError as e: + raise missing_extra('boto3', 'agentcore') from e + kwargs: dict[str, Any] = {} + if self.config.region_name: + kwargs['region_name'] = self.config.region_name + self._client = boto3.client('bedrock-agentcore', **kwargs) + return self._client + + async def _invoke(self, method_name: str, **kwargs: Any) -> Any: + client = self._get_client() + method = getattr(client, method_name) + return await asyncio.to_thread(method, **kwargs) + + def _actor_id(self, app_name: str, user_id: str) -> str: + return f'{app_name}{_ACTOR_SEPARATOR}{user_id}' + + def _user_id_from_actor(self, actor_id: str, app_name: str) -> Optional[str]: + prefix = f'{app_name}{_ACTOR_SEPARATOR}' + if not actor_id.startswith(prefix): + return None + return actor_id[len(prefix) :] + + async def create_session( + self, + *, + app_name: str, + user_id: str, + state: Optional[dict[str, Any]] = None, + session_id: Optional[str] = None, + ) -> Session: + """Creates a new session. + + AgentCore has no create-session API, so this writes a bootstrap event + (extraction skipped) so later ``get_session`` / ``list_sessions`` see it. + """ + sid = session_id or platform_uuid.new_uuid() + existing = await self._list_all_events( + app_name=app_name, + user_id=user_id, + session_id=sid, + include_payloads=False, + ) + if existing: + raise AlreadyExistsError( + f'Session {sid} already exists for user {user_id} in app {app_name}.' + ) + + initial_state = dict(state or {}) + persisted_state = { + k: v + for k, v in initial_state.items() + if not k.startswith(State.TEMP_PREFIX) + } + now = datetime.now(timezone.utc) + await self._invoke( + 'create_event', + memoryId=self.config.memory_id, + actorId=self._actor_id(app_name, user_id), + sessionId=sid, + eventTimestamp=now, + extractionMode='SKIP', + payload=[{ + 'blob': json.dumps({ + 'adk_kind': _BOOTSTRAP_KIND, + 'state': persisted_state, + }) + }], + ) + return Session( + id=sid, + app_name=app_name, + user_id=user_id, + state=initial_state, + events=[], + last_update_time=now.timestamp(), + ) + + async def get_session( + self, + *, + app_name: str, + user_id: str, + session_id: str, + config: Optional[GetSessionConfig] = None, + ) -> Optional[Session]: + """Gets a session and reconstructs ADK events from AgentCore payloads.""" + raw_events = await self._list_all_events( + app_name=app_name, + user_id=user_id, + session_id=session_id, + include_payloads=True, + ) + if not raw_events: + return None + + raw_events.sort(key=lambda e: _unix_timestamp(e.get('eventTimestamp'))) + + bootstrap_state: dict[str, Any] = {} + events: list[Event] = [] + last_update = _unix_timestamp(raw_events[-1].get('eventTimestamp')) + for raw in raw_events: + kind, payload_state, adk_event = self._parse_payload(raw) + if kind == _BOOTSTRAP_KIND: + bootstrap_state = dict(payload_state or {}) + continue + if adk_event is not None: + events.append(adk_event) + + session_state = dict(bootstrap_state) + for event in events: + if event.actions and event.actions.state_delta: + for key, value in event.actions.state_delta.items(): + if not key.startswith(State.TEMP_PREFIX): + session_state[key] = value + + if config is not None: + if config.after_timestamp is not None: + events = [e for e in events if e.timestamp >= config.after_timestamp] + if config.num_recent_events is not None: + events = ( + events[-config.num_recent_events :] + if config.num_recent_events + else [] + ) + + return Session( + id=session_id, + app_name=app_name, + user_id=user_id, + state=session_state, + events=events, + last_update_time=last_update, + ) + + async def list_sessions( + self, *, app_name: str, user_id: Optional[str] = None + ) -> ListSessionsResponse: + """Lists sessions for a user, or all users of the app when user_id is None.""" + if user_id is not None: + actor_ids = [self._actor_id(app_name, user_id)] + else: + actor_ids = await self._list_actor_ids(app_name) + + sessions: list[Session] = [] + for actor_id in actor_ids: + listed_user = ( + user_id + if user_id is not None + else self._user_id_from_actor(actor_id, app_name) + ) + if listed_user is None: + continue + summaries = await self._list_session_summaries(actor_id) + for item in summaries: + created = item.get('createdAt') + sessions.append( + Session( + id=item['sessionId'], + app_name=app_name, + user_id=listed_user, + events=[], + last_update_time=_unix_timestamp(created), + ) + ) + + sessions.sort(key=lambda s: s.last_update_time) + return ListSessionsResponse(sessions=sessions) + + async def delete_session( + self, *, app_name: str, user_id: str, session_id: str + ) -> None: + """Deletes a session by deleting every AgentCore event in it.""" + raw_events = await self._list_all_events( + app_name=app_name, + user_id=user_id, + session_id=session_id, + include_payloads=False, + ) + actor_id = self._actor_id(app_name, user_id) + for raw in raw_events: + event_id = raw.get('eventId') + if not event_id: + continue + await self._invoke( + 'delete_event', + memoryId=self.config.memory_id, + actorId=actor_id, + sessionId=session_id, + eventId=event_id, + ) + + async def append_event(self, session: Session, event: Event) -> Event: + """Appends an event locally, then flushes it to AgentCore STM.""" + if event.partial: + return event + event = await super().append_event(session, event) + payload: list[dict[str, Any]] = [{'blob': event.model_dump_json()}] + text = self._extract_text(event) + if text: + payload.insert( + 0, + { + 'conversational': { + 'role': self._conversational_role(event), + 'content': {'text': text}, + } + }, + ) + timestamp = datetime.fromtimestamp(event.timestamp, tz=timezone.utc) + await self._invoke( + 'create_event', + memoryId=self.config.memory_id, + actorId=self._actor_id(session.app_name, session.user_id), + sessionId=session.id, + eventTimestamp=timestamp, + payload=payload, + ) + session.last_update_time = event.timestamp + return event + + def _extract_text(self, event: Event) -> Optional[str]: + if not event.content or not event.content.parts: + return None + texts = [p.text for p in event.content.parts if p.text] + return '\n'.join(texts) if texts else None + + def _conversational_role(self, event: Event) -> str: + # AgentCore only accepts ASSISTANT | USER | TOOL | OTHER (uppercase). + parts = event.content.parts if event.content else None + if parts and any(getattr(p, 'function_response', None) for p in parts): + return 'TOOL' + if event.author == 'user': + return 'USER' + return 'ASSISTANT' + + def _parse_payload( + self, raw: dict[str, Any] + ) -> tuple[Optional[str], Optional[dict[str, Any]], Optional[Event]]: + """Returns (bootstrap kind, bootstrap state, ADK event).""" + for item in raw.get('payload') or []: + blob = item.get('blob') + if not isinstance(blob, str): + continue + try: + data = json.loads(blob) + except json.JSONDecodeError: + continue + if isinstance(data, dict) and data.get('adk_kind') == _BOOTSTRAP_KIND: + state = data.get('state') + return ( + _BOOTSTRAP_KIND, + state if isinstance(state, dict) else {}, + None, + ) + try: + return None, None, Event.model_validate_json(blob) + except Exception: # pylint: disable=broad-except + logger.debug('Skipping unreadable AgentCore event blob.', exc_info=True) + continue + return None, None, None + + async def _list_all_events( + self, + *, + app_name: str, + user_id: str, + session_id: str, + include_payloads: bool = True, + ) -> list[dict[str, Any]]: + events: list[dict[str, Any]] = [] + next_token: Optional[str] = None + actor_id = self._actor_id(app_name, user_id) + while True: + params: dict[str, Any] = { + 'memoryId': self.config.memory_id, + 'actorId': actor_id, + 'sessionId': session_id, + 'includePayloads': include_payloads, + 'maxResults': _PAGE_SIZE, + } + if next_token is not None: + params['nextToken'] = next_token + response = await self._invoke('list_events', **params) + events.extend(response.get('events') or []) + next_token = response.get('nextToken') + if not next_token: + break + return events + + async def _list_session_summaries( + self, actor_id: str + ) -> list[dict[str, Any]]: + summaries: list[dict[str, Any]] = [] + next_token: Optional[str] = None + while True: + params: dict[str, Any] = { + 'memoryId': self.config.memory_id, + 'actorId': actor_id, + 'maxResults': _PAGE_SIZE, + } + if next_token is not None: + params['nextToken'] = next_token + response = await self._invoke('list_sessions', **params) + summaries.extend(response.get('sessionSummaries') or []) + next_token = response.get('nextToken') + if not next_token: + break + return summaries + + async def _list_actor_ids(self, app_name: str) -> list[str]: + actor_ids: list[str] = [] + next_token: Optional[str] = None + prefix = f'{app_name}{_ACTOR_SEPARATOR}' + while True: + params: dict[str, Any] = { + 'memoryId': self.config.memory_id, + 'maxResults': _PAGE_SIZE, + } + if next_token is not None: + params['nextToken'] = next_token + response = await self._invoke('list_actors', **params) + for item in response.get('actorSummaries') or []: + actor_id = item.get('actorId') + if isinstance(actor_id, str) and actor_id.startswith(prefix): + actor_ids.append(actor_id) + next_token = response.get('nextToken') + if not next_token: + break + return actor_ids diff --git a/src/google/adk/integrations/agentcore/_config.py b/src/google/adk/integrations/agentcore/_config.py new file mode 100644 index 00000000000..b3b5507a441 --- /dev/null +++ b/src/google/adk/integrations/agentcore/_config.py @@ -0,0 +1,37 @@ +# Copyright 2026 Google LLC +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +"""Configuration for AgentCore Memory session storage.""" + +from __future__ import annotations + +from typing import Optional + +from pydantic import BaseModel +from pydantic import Field + + +class AgentCoreSessionServiceConfig(BaseModel): + """Configuration for AgentCoreSessionService.""" + + memory_id: str = Field( + description='AgentCore Memory resource id (memoryId).', + ) + region_name: Optional[str] = Field( + default=None, + description=( + 'AWS region for the bedrock-agentcore client. If unset, boto3 uses' + ' the default credential/region chain.' + ), + ) diff --git a/tests/unittests/integrations/agentcore/__init__.py b/tests/unittests/integrations/agentcore/__init__.py new file mode 100644 index 00000000000..307b2052b75 --- /dev/null +++ b/tests/unittests/integrations/agentcore/__init__.py @@ -0,0 +1,15 @@ +# Copyright 2026 Google LLC +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +"""Unit tests for AgentCore integrations.""" diff --git a/tests/unittests/integrations/agentcore/_fake_agentcore.py b/tests/unittests/integrations/agentcore/_fake_agentcore.py new file mode 100644 index 00000000000..bfeaf762364 --- /dev/null +++ b/tests/unittests/integrations/agentcore/_fake_agentcore.py @@ -0,0 +1,85 @@ +# Copyright 2026 Google LLC +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +"""In-memory stand-in for the bedrock-agentcore client methods ADK calls.""" + +from __future__ import annotations + +from datetime import datetime +from datetime import timezone +from typing import Any + + +class FakeAgentCoreClient: + """Stores Actor/session/event records the way AgentCore Memory does.""" + + def __init__(self) -> None: + # (actor_id, session_id) -> list of event dicts + self._events: dict[tuple[str, str], list[dict[str, Any]]] = {} + self._counter = 0 + + def create_event(self, **kwargs: Any) -> dict[str, Any]: + actor_id = kwargs['actorId'] + session_id = kwargs['sessionId'] + self._counter += 1 + timestamp = kwargs.get('eventTimestamp') or datetime.now(timezone.utc) + event = { + 'memoryId': kwargs['memoryId'], + 'actorId': actor_id, + 'sessionId': session_id, + 'eventId': f'evt-{self._counter}', + 'eventTimestamp': timestamp, + 'payload': kwargs.get('payload') or [], + } + if 'extractionMode' in kwargs: + event['extractionMode'] = kwargs['extractionMode'] + key = (actor_id, session_id) + self._events.setdefault(key, []).append(event) + return {'event': dict(event)} + + def list_events(self, **kwargs: Any) -> dict[str, Any]: + key = (kwargs['actorId'], kwargs['sessionId']) + events = list(self._events.get(key, [])) + include_payloads = kwargs.get('includePayloads', True) + if not include_payloads: + events = [ + {k: v for k, v in event.items() if k != 'payload'} for event in events + ] + return {'events': events} + + def list_sessions(self, **kwargs: Any) -> dict[str, Any]: + actor_id = kwargs['actorId'] + summaries: list[dict[str, Any]] = [] + for (stored_actor, session_id), events in self._events.items(): + if stored_actor != actor_id or not events: + continue + summaries.append({ + 'sessionId': session_id, + 'actorId': stored_actor, + 'createdAt': events[0]['eventTimestamp'], + }) + return {'sessionSummaries': summaries} + + def list_actors(self, **kwargs: Any) -> dict[str, Any]: + actor_ids = sorted({actor_id for actor_id, _ in self._events}) + return {'actorSummaries': [{'actorId': a} for a in actor_ids]} + + def delete_event(self, **kwargs: Any) -> dict[str, Any]: + key = (kwargs['actorId'], kwargs['sessionId']) + event_id = kwargs['eventId'] + stored = self._events.get(key, []) + self._events[key] = [e for e in stored if e.get('eventId') != event_id] + if not self._events[key]: + self._events.pop(key, None) + return {'eventId': event_id} diff --git a/tests/unittests/integrations/agentcore/test_agentcore_session_service.py b/tests/unittests/integrations/agentcore/test_agentcore_session_service.py new file mode 100644 index 00000000000..a6e5ac6c489 --- /dev/null +++ b/tests/unittests/integrations/agentcore/test_agentcore_session_service.py @@ -0,0 +1,293 @@ +# Copyright 2026 Google LLC +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +"""Unit tests for AgentCoreSessionService.""" + +from __future__ import annotations + +from google.adk.errors.already_exists_error import AlreadyExistsError +from google.adk.events.event import Event +from google.adk.events.event import EventActions +from google.adk.integrations.agentcore._agentcore_session_service import AgentCoreSessionService +from google.adk.integrations.agentcore._config import AgentCoreSessionServiceConfig +from google.adk.sessions.base_session_service import GetSessionConfig +from google.genai import types +import pytest + +from ._fake_agentcore import FakeAgentCoreClient + + +@pytest.fixture +def fake_client(): + return FakeAgentCoreClient() + + +@pytest.fixture +def session_service(fake_client): + return AgentCoreSessionService( + config=AgentCoreSessionServiceConfig(memory_id='mem-1'), + client=fake_client, + ) + + +def test_memory_id_is_required(): + with pytest.raises(ValueError, match='memory_id is required'): + AgentCoreSessionService() + + +def _user_event(text: str, **kwargs) -> Event: + return Event( + author='user', + content=types.Content(role='user', parts=[types.Part(text=text)]), + **kwargs, + ) + + +@pytest.mark.asyncio +async def test_create_session(session_service): + session = await session_service.create_session( + app_name='app1', + user_id='user1', + state={'key1': 'val1', 'temp:scratch': 'nope'}, + ) + + assert session.app_name == 'app1' + assert session.user_id == 'user1' + assert session.state['key1'] == 'val1' + assert session.state['temp:scratch'] == 'nope' + assert session.id is not None + + fetched = await session_service.get_session( + app_name='app1', user_id='user1', session_id=session.id + ) + assert fetched is not None + assert fetched.state['key1'] == 'val1' + assert 'temp:scratch' not in fetched.state + + +@pytest.mark.asyncio +async def test_create_session_already_exists(session_service): + await session_service.create_session( + app_name='app1', user_id='user1', session_id='sess_123' + ) + + with pytest.raises(AlreadyExistsError): + await session_service.create_session( + app_name='app1', user_id='user1', session_id='sess_123' + ) + + +@pytest.mark.asyncio +async def test_get_session_not_found(session_service): + fetched = await session_service.get_session( + app_name='app1', user_id='user1', session_id='missing' + ) + assert fetched is None + + +@pytest.mark.asyncio +async def test_append_event_round_trips_blob_and_text( + session_service, fake_client +): + session = await session_service.create_session( + app_name='app1', user_id='user1', session_id='s1' + ) + event = _user_event('hello there') + await session_service.append_event(session, event) + + fetched = await session_service.get_session( + app_name='app1', user_id='user1', session_id='s1' + ) + assert fetched is not None + assert len(fetched.events) == 1 + assert fetched.events[0].author == 'user' + assert fetched.events[0].content.parts[0].text == 'hello there' + + stored = fake_client._events[('app1:user1', 's1')] + payloads = stored[-1]['payload'] + assert payloads[0]['conversational']['role'] == 'USER' + assert payloads[0]['conversational']['content']['text'] == 'hello there' + assert 'blob' in payloads[1] + + assistant = Event( + author='simple_agent', + content=types.Content( + role='model', parts=[types.Part(text='hi from the model')] + ), + ) + await session_service.append_event(fetched, assistant) + assistant_payload = fake_client._events[('app1:user1', 's1')][-1]['payload'] + assert assistant_payload[0]['conversational']['role'] == 'ASSISTANT' + + +@pytest.mark.asyncio +async def test_append_event_tool_role(session_service, fake_client): + session = await session_service.create_session( + app_name='app1', user_id='user1' + ) + event = Event( + author='tool', + content=types.Content( + role='user', + parts=[ + types.Part( + function_response=types.FunctionResponse( + name='lookup', response={'ok': True} + ) + ) + ], + ), + ) + # No text → conversational payload is omitted; role helper still maps TOOL. + await session_service.append_event(session, event) + assert session_service._conversational_role(event) == 'TOOL' + + +@pytest.mark.asyncio +async def test_partial_events_are_not_written(session_service, fake_client): + session = await session_service.create_session( + app_name='app1', user_id='user1', session_id='s1' + ) + await session_service.append_event( + session, Event(author='user', partial=True) + ) + + fetched = await session_service.get_session( + app_name='app1', user_id='user1', session_id='s1' + ) + assert fetched is not None + assert fetched.events == [] + # Bootstrap only. + assert len(fake_client._events[('app1:user1', 's1')]) == 1 + + +@pytest.mark.asyncio +async def test_get_session_with_event_filter(session_service): + session = await session_service.create_session( + app_name='app1', user_id='user1' + ) + for i in range(5): + await session_service.append_event( + session, Event(author=f'user_{i}', timestamp=float(100 + i)) + ) + + fetched = await session_service.get_session( + app_name='app1', + user_id='user1', + session_id=session.id, + config=GetSessionConfig(num_recent_events=2), + ) + assert fetched is not None + assert len(fetched.events) == 2 + assert fetched.events[-1].author == 'user_4' + + fetched_zero = await session_service.get_session( + app_name='app1', + user_id='user1', + session_id=session.id, + config=GetSessionConfig(num_recent_events=0), + ) + assert fetched_zero is not None + assert fetched_zero.events == [] + + fetched_after = await session_service.get_session( + app_name='app1', + user_id='user1', + session_id=session.id, + config=GetSessionConfig(after_timestamp=103.0), + ) + assert fetched_after is not None + assert [e.author for e in fetched_after.events] == ['user_3', 'user_4'] + + +@pytest.mark.asyncio +async def test_append_event_and_state_delta(session_service): + session = await session_service.create_session(app_name='app1', user_id='u1') + event = Event( + author='agent', + actions=EventActions( + state_delta={ + 'count': 1, + 'user:score': 100, + 'app:status': 'active', + 'temp:scratch': 'gone', + } + ), + ) + await session_service.append_event(session, event) + + fetched = await session_service.get_session( + app_name='app1', user_id='u1', session_id=session.id + ) + assert fetched is not None + assert len(fetched.events) == 1 + assert fetched.state['count'] == 1 + assert fetched.state['user:score'] == 100 + assert fetched.state['app:status'] == 'active' + assert 'temp:scratch' not in fetched.state + + +@pytest.mark.asyncio +async def test_list_sessions(session_service): + await session_service.create_session( + app_name='app1', user_id='u1', session_id='s1' + ) + await session_service.create_session( + app_name='app1', user_id='u1', session_id='s2' + ) + await session_service.create_session( + app_name='app1', user_id='u2', session_id='s3' + ) + await session_service.create_session( + app_name='app2', user_id='u1', session_id='s1' + ) + + resp_u1 = await session_service.list_sessions(app_name='app1', user_id='u1') + assert {s.id for s in resp_u1.sessions} == {'s1', 's2'} + assert all(s.events == [] for s in resp_u1.sessions) + + resp_all = await session_service.list_sessions(app_name='app1') + assert {s.id for s in resp_all.sessions} == {'s1', 's2', 's3'} + assert {s.user_id for s in resp_all.sessions} == {'u1', 'u2'} + + +@pytest.mark.asyncio +async def test_delete_session(session_service): + await session_service.create_session( + app_name='app1', user_id='u1', session_id='to_delete' + ) + session = await session_service.get_session( + app_name='app1', user_id='u1', session_id='to_delete' + ) + assert session is not None + await session_service.append_event(session, _user_event('bye')) + + await session_service.delete_session( + app_name='app1', user_id='u1', session_id='to_delete' + ) + fetched = await session_service.get_session( + app_name='app1', user_id='u1', session_id='to_delete' + ) + assert fetched is None + + +@pytest.mark.asyncio +async def test_apps_do_not_share_actor_namespace(session_service): + await session_service.create_session( + app_name='app1', user_id='u1', session_id='shared-id' + ) + other = await session_service.get_session( + app_name='app2', user_id='u1', session_id='shared-id' + ) + assert other is None