From f8ed4ae9ef7ff1f1d1f08a623002c769ceeb3bc3 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Gon=C3=A7alo=20Faustino?= Date: Wed, 8 Apr 2026 15:46:25 +0100 Subject: [PATCH 1/3] feat(adk): add native Bedrock embedding support for agent memory MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Add _embed_bedrock() to KagentMemoryService so that memory works with AWS Bedrock embedding models (e.g. amazon.titan-embed-text-v1/v2) without falling back to the OpenAI-compatible client. The implementation uses boto3's invoke_model API, consistent with how KAgentBedrockLlm already handles LLM calls. Region is resolved from AWS_DEFAULT_REGION / AWS_REGION env vars (same credential chain as the Bedrock LLM provider). Each text is embedded individually since the Titan Embedding API accepts a single inputText per invocation; asyncio.gather parallelises the calls. The upstream truncation-to-768 + L2-normalization logic handles dimension differences across Titan v1 (1536d) and v2 (1024d) models, so no model-specific dimension parameter is needed. Signed-off-by: Gonçalo Faustino Made-with: Cursor --- .../src/kagent/adk/_memory_service.py | 35 +++++++++++++++++ .../tests/unittests/test_embedding.py | 38 +++++++++++++++++++ 2 files changed, 73 insertions(+) diff --git a/python/packages/kagent-adk/src/kagent/adk/_memory_service.py b/python/packages/kagent-adk/src/kagent/adk/_memory_service.py index 797c9ef3f..b5fdcef89 100644 --- a/python/packages/kagent-adk/src/kagent/adk/_memory_service.py +++ b/python/packages/kagent-adk/src/kagent/adk/_memory_service.py @@ -368,6 +368,8 @@ async def _call_embedding_provider( return await self._embed_ollama(model_name, texts, api_base) if provider in ("vertex_ai", "gemini"): return await self._embed_google(provider, model_name, texts) + if provider == "bedrock": + return await self._embed_bedrock(model_name, texts) # Unknown provider — try OpenAI-compatible as a fallback logger.warning("Unknown embedding provider '%s'; attempting OpenAI-compatible call.", provider) return await self._embed_openai("openai", model_name, texts, api_base) @@ -437,6 +439,39 @@ async def _embed_google( ) return [list(emb.values) for emb in response.embeddings] + async def _embed_bedrock( + self, + model_name: str, + texts: List[str], + ) -> List[List[float]]: + """Embed using the AWS Bedrock Titan Embedding API via boto3. + + Uses the same credential chain (env vars, IRSA, instance profile) as + KAgentBedrockLlm. Each text is embedded individually because the + Titan Embedding API accepts a single ``inputText`` per invocation. + """ + import os + + import boto3 + + region = os.environ.get("AWS_DEFAULT_REGION") or os.environ.get("AWS_REGION") or "us-east-1" + client = boto3.client("bedrock-runtime", region_name=region) + + async def _invoke_single(text: str) -> List[float]: + body = json.dumps({"inputText": text}) + response = await asyncio.to_thread( + client.invoke_model, + modelId=model_name, + body=body, + contentType="application/json", + accept="application/json", + ) + result = json.loads(response["body"].read()) + return result["embedding"] + + embeddings = await asyncio.gather(*[_invoke_single(t) for t in texts]) + return list(embeddings) + async def _summarize_session_content_async( self, content: str, diff --git a/python/packages/kagent-adk/tests/unittests/test_embedding.py b/python/packages/kagent-adk/tests/unittests/test_embedding.py index 54464c253..fa15c3188 100644 --- a/python/packages/kagent-adk/tests/unittests/test_embedding.py +++ b/python/packages/kagent-adk/tests/unittests/test_embedding.py @@ -150,3 +150,41 @@ async def test_embedding_shorter_than_768_rejected(self): mock_cls.return_value = instance result = await svc._generate_embedding_async("test") assert result == [] + + @pytest.mark.asyncio + async def test_bedrock_embed(self): + svc = make_service(provider="bedrock", model="amazon.titan-embed-text-v1") + vec = [0.1] * 1536 + mock_response = mock.MagicMock() + mock_response.__getitem__ = lambda self, key: {"body": mock.MagicMock(read=lambda: b'{"embedding": ' + str(vec).encode() + b"}")}[key] + mock_client = mock.MagicMock() + mock_client.invoke_model = mock.MagicMock(return_value={"body": mock.MagicMock(read=lambda: b'{"embedding": ' + str(vec).encode() + b"}")}) + + with mock.patch("boto3.client", return_value=mock_client): + result = await svc._generate_embedding_async("hello world") + assert len(result) == 768 + mock_client.invoke_model.assert_called_once() + + @pytest.mark.asyncio + async def test_bedrock_embed_uses_region_from_env(self): + svc = make_service(provider="bedrock", model="amazon.titan-embed-text-v1") + vec = [0.5] * 1536 + mock_client = mock.MagicMock() + mock_client.invoke_model = mock.MagicMock(return_value={"body": mock.MagicMock(read=lambda: b'{"embedding": ' + str(vec).encode() + b"}")}) + + with ( + mock.patch.dict("os.environ", {"AWS_REGION": "eu-west-1"}), + mock.patch("boto3.client", return_value=mock_client) as mock_boto, + ): + await svc._generate_embedding_async("test") + mock_boto.assert_called_once_with("bedrock-runtime", region_name="eu-west-1") + + @pytest.mark.asyncio + async def test_bedrock_embed_error_returns_empty(self): + svc = make_service(provider="bedrock", model="amazon.titan-embed-text-v1") + mock_client = mock.MagicMock() + mock_client.invoke_model = mock.MagicMock(side_effect=Exception("Bedrock API error")) + + with mock.patch("boto3.client", return_value=mock_client): + result = await svc._generate_embedding_async("test") + assert result == [] From 32538d805a9bb307ad1fdc46b4aaef24cb042466 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Gon=C3=A7alo=20Faustino?= Date: Wed, 8 Apr 2026 15:55:50 +0100 Subject: [PATCH 2/3] fix(adk): address review feedback on Bedrock embedding tests MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Remove unused mock_response in test_bedrock_embed - Override AWS_DEFAULT_REGION in region env test to prevent flakiness when the runner environment has it set Signed-off-by: Gonçalo Faustino Made-with: Cursor --- python/packages/kagent-adk/tests/unittests/test_embedding.py | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/python/packages/kagent-adk/tests/unittests/test_embedding.py b/python/packages/kagent-adk/tests/unittests/test_embedding.py index fa15c3188..94e9fa546 100644 --- a/python/packages/kagent-adk/tests/unittests/test_embedding.py +++ b/python/packages/kagent-adk/tests/unittests/test_embedding.py @@ -155,8 +155,6 @@ async def test_embedding_shorter_than_768_rejected(self): async def test_bedrock_embed(self): svc = make_service(provider="bedrock", model="amazon.titan-embed-text-v1") vec = [0.1] * 1536 - mock_response = mock.MagicMock() - mock_response.__getitem__ = lambda self, key: {"body": mock.MagicMock(read=lambda: b'{"embedding": ' + str(vec).encode() + b"}")}[key] mock_client = mock.MagicMock() mock_client.invoke_model = mock.MagicMock(return_value={"body": mock.MagicMock(read=lambda: b'{"embedding": ' + str(vec).encode() + b"}")}) @@ -173,7 +171,7 @@ async def test_bedrock_embed_uses_region_from_env(self): mock_client.invoke_model = mock.MagicMock(return_value={"body": mock.MagicMock(read=lambda: b'{"embedding": ' + str(vec).encode() + b"}")}) with ( - mock.patch.dict("os.environ", {"AWS_REGION": "eu-west-1"}), + mock.patch.dict("os.environ", {"AWS_REGION": "eu-west-1", "AWS_DEFAULT_REGION": ""}), mock.patch("boto3.client", return_value=mock_client) as mock_boto, ): await svc._generate_embedding_async("test") From 6bd89dece5b633a84af7aec1d8b46a8bc4376a22 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Gon=C3=A7alo=20Faustino?= Date: Thu, 9 Apr 2026 17:24:17 +0100 Subject: [PATCH 3/3] style(adk): fix ruff format on Bedrock embedding tests MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Wrap long mock_client.invoke_model lines to comply with ruff line-length formatting rules. Signed-off-by: Gonçalo Faustino Made-with: Cursor --- .../packages/kagent-adk/tests/unittests/test_embedding.py | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) diff --git a/python/packages/kagent-adk/tests/unittests/test_embedding.py b/python/packages/kagent-adk/tests/unittests/test_embedding.py index 94e9fa546..ffeffd78a 100644 --- a/python/packages/kagent-adk/tests/unittests/test_embedding.py +++ b/python/packages/kagent-adk/tests/unittests/test_embedding.py @@ -156,7 +156,9 @@ async def test_bedrock_embed(self): svc = make_service(provider="bedrock", model="amazon.titan-embed-text-v1") vec = [0.1] * 1536 mock_client = mock.MagicMock() - mock_client.invoke_model = mock.MagicMock(return_value={"body": mock.MagicMock(read=lambda: b'{"embedding": ' + str(vec).encode() + b"}")}) + mock_client.invoke_model = mock.MagicMock( + return_value={"body": mock.MagicMock(read=lambda: b'{"embedding": ' + str(vec).encode() + b"}")} + ) with mock.patch("boto3.client", return_value=mock_client): result = await svc._generate_embedding_async("hello world") @@ -168,7 +170,9 @@ async def test_bedrock_embed_uses_region_from_env(self): svc = make_service(provider="bedrock", model="amazon.titan-embed-text-v1") vec = [0.5] * 1536 mock_client = mock.MagicMock() - mock_client.invoke_model = mock.MagicMock(return_value={"body": mock.MagicMock(read=lambda: b'{"embedding": ' + str(vec).encode() + b"}")}) + mock_client.invoke_model = mock.MagicMock( + return_value={"body": mock.MagicMock(read=lambda: b'{"embedding": ' + str(vec).encode() + b"}")} + ) with ( mock.patch.dict("os.environ", {"AWS_REGION": "eu-west-1", "AWS_DEFAULT_REGION": ""}),