From 47c0a0ac97151257f2085e01a68ccf2f6b41f0e1 Mon Sep 17 00:00:00 2001 From: Yufeng He <40085740+he-yufeng@users.noreply.github.com> Date: Tue, 4 Aug 2026 05:38:27 +0800 Subject: [PATCH 1/4] fix(redis): scope RedisHistoryProvider keys by source_id Two providers with different source_ids but the same key_prefix shared one Redis list per session, so a write-only audit sink contaminated the primary provider's loaded history, and clear() on one deleted the other's conversation. The key now includes source_id, matching the Cosmos provider's scoping. Existing keys written under the old layout are left in place; deleting them would risk removing a sibling provider's data, and they simply become unreadable by the new code. --- .../_history_provider.py | 2 +- python/packages/redis/tests/test_providers.py | 27 ++++++++++++++++--- 2 files changed, 24 insertions(+), 5 deletions(-) diff --git a/python/packages/redis/agent_framework_redis/_history_provider.py b/python/packages/redis/agent_framework_redis/_history_provider.py index a7703db8a2..34305413ba 100644 --- a/python/packages/redis/agent_framework_redis/_history_provider.py +++ b/python/packages/redis/agent_framework_redis/_history_provider.py @@ -108,7 +108,7 @@ def __init__( def _redis_key(self, session_id: str | None) -> str: """Get the Redis key for a given session's messages.""" - return f"{self.key_prefix}:{session_id or 'default'}" + return f"{self.key_prefix}:{self.source_id}:{session_id or 'default'}" async def get_messages( self, diff --git a/python/packages/redis/tests/test_providers.py b/python/packages/redis/tests/test_providers.py index 55aee29662..ca09d53e6e 100644 --- a/python/packages/redis/tests/test_providers.py +++ b/python/packages/redis/tests/test_providers.py @@ -420,8 +420,16 @@ def test_key_format(self, mock_redis_client: MagicMock): mock_from_url.return_value = mock_redis_client provider = RedisHistoryProvider("mem", redis_url="redis://localhost:6379", key_prefix="msgs") - assert provider._redis_key("session-123") == "msgs:session-123" - assert provider._redis_key(None) == "msgs:default" + assert provider._redis_key("session-123") == "msgs:mem:session-123" + assert provider._redis_key(None) == "msgs:mem:default" + + def test_keys_isolated_per_source_id(self, mock_redis_client: MagicMock): + with patch("agent_framework_redis._history_provider.redis.from_url") as mock_from_url: + mock_from_url.return_value = mock_redis_client + first = RedisHistoryProvider("audit", redis_url="redis://localhost:6379", key_prefix="msgs") + second = RedisHistoryProvider("primary", redis_url="redis://localhost:6379", key_prefix="msgs") + + assert first._redis_key("s1") != second._redis_key("s1") class TestRedisHistoryProviderGetMessages: @@ -482,7 +490,7 @@ async def test_max_messages_trimming(self, mock_redis_client: MagicMock): await provider.save_messages("s1", [Message(role="user", contents=["msg"])]) - mock_redis_client.ltrim.assert_called_once_with("chat_messages:s1", -10, -1) + mock_redis_client.ltrim.assert_called_once_with("chat_messages:mem:s1", -10, -1) async def test_no_trim_when_under_limit(self, mock_redis_client: MagicMock): mock_redis_client.llen = AsyncMock(return_value=3) @@ -503,7 +511,18 @@ async def test_clear_calls_delete(self, mock_redis_client: MagicMock): provider = RedisHistoryProvider("mem", redis_url="redis://localhost:6379") await provider.clear("session-1") - mock_redis_client.delete.assert_called_once_with("chat_messages:session-1") + mock_redis_client.delete.assert_called_once_with("chat_messages:mem:session-1") + + async def test_clear_leaves_other_source_ids_untouched(self, mock_redis_client: MagicMock): + with patch("agent_framework_redis._history_provider.redis.from_url") as mock_from_url: + mock_from_url.return_value = mock_redis_client + audit = RedisHistoryProvider("audit", redis_url="redis://localhost:6379") + primary = RedisHistoryProvider("primary", redis_url="redis://localhost:6379") + + await audit.clear("session-1") + # the destructive case from #7471: clearing one provider must not + # delete the shared session's messages belonging to another provider + mock_redis_client.delete.assert_called_once_with("chat_messages:audit:session-1") class TestRedisHistoryProviderBeforeAfterRun: From d309e32691ce2e597bc6c0b3898fd5d008d1076b Mon Sep 17 00:00:00 2001 From: Yufeng He <40085740+he-yufeng@users.noreply.github.com> Date: Tue, 4 Aug 2026 09:29:22 +0800 Subject: [PATCH 2/4] fix(redis): make the provider key separator collision-proof and use it in tests Colon-joined keys were ambiguous for source ids or session ids containing a colon (a:b + c vs a + b:c). Join with the ASCII unit separator instead. The clear-isolation test now asserts on the other provider's key too, so the unused variable lint is gone as well. --- .../redis/agent_framework_redis/_history_provider.py | 6 +++++- python/packages/redis/tests/test_providers.py | 11 ++++++----- 2 files changed, 11 insertions(+), 6 deletions(-) diff --git a/python/packages/redis/agent_framework_redis/_history_provider.py b/python/packages/redis/agent_framework_redis/_history_provider.py index 34305413ba..f260e3f2da 100644 --- a/python/packages/redis/agent_framework_redis/_history_provider.py +++ b/python/packages/redis/agent_framework_redis/_history_provider.py @@ -106,9 +106,13 @@ def __init__( else: self._redis_client = redis.from_url(redis_url, decode_responses=True) # type: ignore[no-untyped-call] + # Unit separator: source ids and session ids are opaque strings and can + # legitimately contain ':', which would make colon-joined keys ambiguous. + _KEY_SEP = "\x1f" + def _redis_key(self, session_id: str | None) -> str: """Get the Redis key for a given session's messages.""" - return f"{self.key_prefix}:{self.source_id}:{session_id or 'default'}" + return self._KEY_SEP.join([self.key_prefix, self.source_id, session_id or "default"]) async def get_messages( self, diff --git a/python/packages/redis/tests/test_providers.py b/python/packages/redis/tests/test_providers.py index ca09d53e6e..a73410182d 100644 --- a/python/packages/redis/tests/test_providers.py +++ b/python/packages/redis/tests/test_providers.py @@ -420,8 +420,8 @@ def test_key_format(self, mock_redis_client: MagicMock): mock_from_url.return_value = mock_redis_client provider = RedisHistoryProvider("mem", redis_url="redis://localhost:6379", key_prefix="msgs") - assert provider._redis_key("session-123") == "msgs:mem:session-123" - assert provider._redis_key(None) == "msgs:mem:default" + assert provider._redis_key("session-123") == "msgs\x1fmem\x1fsession-123" + assert provider._redis_key(None) == "msgs\x1fmem\x1fdefault" def test_keys_isolated_per_source_id(self, mock_redis_client: MagicMock): with patch("agent_framework_redis._history_provider.redis.from_url") as mock_from_url: @@ -490,7 +490,7 @@ async def test_max_messages_trimming(self, mock_redis_client: MagicMock): await provider.save_messages("s1", [Message(role="user", contents=["msg"])]) - mock_redis_client.ltrim.assert_called_once_with("chat_messages:mem:s1", -10, -1) + mock_redis_client.ltrim.assert_called_once_with("chat_messages\x1fmem\x1fs1", -10, -1) async def test_no_trim_when_under_limit(self, mock_redis_client: MagicMock): mock_redis_client.llen = AsyncMock(return_value=3) @@ -511,7 +511,7 @@ async def test_clear_calls_delete(self, mock_redis_client: MagicMock): provider = RedisHistoryProvider("mem", redis_url="redis://localhost:6379") await provider.clear("session-1") - mock_redis_client.delete.assert_called_once_with("chat_messages:mem:session-1") + mock_redis_client.delete.assert_called_once_with("chat_messages\x1fmem\x1fsession-1") async def test_clear_leaves_other_source_ids_untouched(self, mock_redis_client: MagicMock): with patch("agent_framework_redis._history_provider.redis.from_url") as mock_from_url: @@ -522,7 +522,8 @@ async def test_clear_leaves_other_source_ids_untouched(self, mock_redis_client: await audit.clear("session-1") # the destructive case from #7471: clearing one provider must not # delete the shared session's messages belonging to another provider - mock_redis_client.delete.assert_called_once_with("chat_messages:audit:session-1") + mock_redis_client.delete.assert_called_once_with("chat_messages\x1faudit\x1fsession-1") + assert primary._redis_key("session-1") not in mock_redis_client.delete.call_args.args class TestRedisHistoryProviderBeforeAfterRun: From 22a5adff3892451500badd9e45193b093db664c4 Mon Sep 17 00:00:00 2001 From: Yufeng He <40085740+he-yufeng@users.noreply.github.com> Date: Sun, 16 Aug 2026 09:35:19 +0900 Subject: [PATCH 3/4] fix(redis): length-prefix history key components and migrate pre-scoping keys lazily --- .../_history_provider.py | 28 ++++++++-- python/packages/redis/tests/test_providers.py | 53 +++++++++++++++++-- 2 files changed, 71 insertions(+), 10 deletions(-) diff --git a/python/packages/redis/agent_framework_redis/_history_provider.py b/python/packages/redis/agent_framework_redis/_history_provider.py index f260e3f2da..25cf0ac8c9 100644 --- a/python/packages/redis/agent_framework_redis/_history_provider.py +++ b/python/packages/redis/agent_framework_redis/_history_provider.py @@ -9,6 +9,7 @@ from __future__ import annotations from collections.abc import Sequence +from contextlib import suppress from typing import Any, ClassVar import redis.asyncio as redis @@ -106,13 +107,20 @@ def __init__( else: self._redis_client = redis.from_url(redis_url, decode_responses=True) # type: ignore[no-untyped-call] - # Unit separator: source ids and session ids are opaque strings and can - # legitimately contain ':', which would make colon-joined keys ambiguous. - _KEY_SEP = "\x1f" + # Keys length-prefix each component (":") so the join stays + # injective no matter which bytes the source/session ids carry; any fixed + # separator can be smuggled inside an opaque id and collide two sessions. + # Sessions written before source_id scoping live under + # ":" and migrate lazily on first read. def _redis_key(self, session_id: str | None) -> str: """Get the Redis key for a given session's messages.""" - return self._KEY_SEP.join([self.key_prefix, self.source_id, session_id or "default"]) + parts = (self.key_prefix, self.source_id, session_id or "default") + return "".join(f"{len(part)}:{part}" for part in parts) + + def _legacy_redis_key(self, session_id: str | None) -> str: + """Pre-scoping key layout, read only to migrate existing sessions.""" + return f"{self.key_prefix}:{session_id or 'default'}" async def get_messages( self, @@ -134,6 +142,16 @@ async def get_messages( mark_feature_used(FeatureIndex.REDIS) key = self._redis_key(session_id) redis_messages: list[str] = await self._redis_client.lrange(key, 0, -1) # type: ignore[misc] + if not redis_messages: + # Lazy migration: a session last written with the pre-scoping key + # layout moves under its new key on first read. renamenx keeps this + # atomic and no-ops if a concurrent write already landed there; the + # legacy list then stays put and remains clearable via clear(). + legacy_key = self._legacy_redis_key(session_id) + if legacy_key != key and await self._redis_client.exists(legacy_key): # type: ignore[misc] + with suppress(Exception): # a legacy key that vanished mid-read is a no-op + await self._redis_client.renamenx(legacy_key, key) # type: ignore[misc] + redis_messages = await self._redis_client.lrange(key, 0, -1) # type: ignore[misc] messages: list[Message] = [] if redis_messages: for serialized in redis_messages: # type: ignore[union-attr] @@ -193,7 +211,7 @@ async def clear(self, session_id: str | None) -> None: Args: session_id: The session ID to clear messages for. """ - await self._redis_client.delete(self._redis_key(session_id)) + await self._redis_client.delete(self._redis_key(session_id), self._legacy_redis_key(session_id)) async def aclose(self) -> None: """Close the Redis connection.""" diff --git a/python/packages/redis/tests/test_providers.py b/python/packages/redis/tests/test_providers.py index a73410182d..381cbed5f2 100644 --- a/python/packages/redis/tests/test_providers.py +++ b/python/packages/redis/tests/test_providers.py @@ -63,6 +63,8 @@ def mock_redis_client(): client.llen = AsyncMock(return_value=0) client.ltrim = AsyncMock() client.delete = AsyncMock() + client.exists = AsyncMock(return_value=0) + client.renamenx = AsyncMock(return_value=0) mock_pipeline = AsyncMock() mock_pipeline.rpush = AsyncMock() @@ -420,8 +422,18 @@ def test_key_format(self, mock_redis_client: MagicMock): mock_from_url.return_value = mock_redis_client provider = RedisHistoryProvider("mem", redis_url="redis://localhost:6379", key_prefix="msgs") - assert provider._redis_key("session-123") == "msgs\x1fmem\x1fsession-123" - assert provider._redis_key(None) == "msgs\x1fmem\x1fdefault" + assert provider._redis_key("session-123") == "4:msgs3:mem11:session-123" + assert provider._redis_key(None) == "4:msgs3:mem7:default" + + def test_key_join_is_injective(self, mock_redis_client: MagicMock): + # moonbox3's review case: any fixed separator can be smuggled inside an + # opaque id, so the components are length-prefixed instead. + with patch("agent_framework_redis._history_provider.redis.from_url") as mock_from_url: + mock_from_url.return_value = mock_redis_client + first = RedisHistoryProvider("audit\x1fx", redis_url="redis://localhost:6379", key_prefix="msgs") + second = RedisHistoryProvider("audit", redis_url="redis://localhost:6379", key_prefix="msgs") + + assert first._redis_key("y") != second._redis_key("x\x1fy") def test_keys_isolated_per_source_id(self, mock_redis_client: MagicMock): with patch("agent_framework_redis._history_provider.redis.from_url") as mock_from_url: @@ -459,6 +471,35 @@ async def test_empty_returns_empty(self, mock_redis_client: MagicMock): messages = await provider.get_messages("s1") assert messages == [] + async def test_migrates_legacy_key_on_first_read(self, mock_redis_client: MagicMock): + msg = Message(role="user", contents=["legacy hello"]) + legacy_payload = json.dumps(msg.to_dict()) + # new key empty, legacy key still holds the pre-scoping data + mock_redis_client.lrange = AsyncMock(side_effect=[[], [legacy_payload]]) + mock_redis_client.exists = AsyncMock(return_value=1) + mock_redis_client.renamenx = AsyncMock(return_value=1) + + with patch("agent_framework_redis._history_provider.redis.from_url") as mock_from_url: + mock_from_url.return_value = mock_redis_client + provider = RedisHistoryProvider("mem", redis_url="redis://localhost:6379") + + messages = await provider.get_messages("s1") + mock_redis_client.renamenx.assert_called_once_with("chat_messages:s1", "13:chat_messages3:mem2:s1") + assert len(messages) == 1 + assert messages[0].text == "legacy hello" + + async def test_no_migration_when_new_key_has_data(self, mock_redis_client: MagicMock): + msg = Message(role="user", contents=["current"]) + mock_redis_client.lrange = AsyncMock(return_value=[json.dumps(msg.to_dict())]) + + with patch("agent_framework_redis._history_provider.redis.from_url") as mock_from_url: + mock_from_url.return_value = mock_redis_client + provider = RedisHistoryProvider("mem", redis_url="redis://localhost:6379") + + messages = await provider.get_messages("s1") + mock_redis_client.exists.assert_not_called() + assert len(messages) == 1 + class TestRedisHistoryProviderSaveMessages: async def test_saves_serialized_messages(self, mock_redis_client: MagicMock): @@ -490,7 +531,7 @@ async def test_max_messages_trimming(self, mock_redis_client: MagicMock): await provider.save_messages("s1", [Message(role="user", contents=["msg"])]) - mock_redis_client.ltrim.assert_called_once_with("chat_messages\x1fmem\x1fs1", -10, -1) + mock_redis_client.ltrim.assert_called_once_with("13:chat_messages3:mem2:s1", -10, -1) async def test_no_trim_when_under_limit(self, mock_redis_client: MagicMock): mock_redis_client.llen = AsyncMock(return_value=3) @@ -511,7 +552,7 @@ async def test_clear_calls_delete(self, mock_redis_client: MagicMock): provider = RedisHistoryProvider("mem", redis_url="redis://localhost:6379") await provider.clear("session-1") - mock_redis_client.delete.assert_called_once_with("chat_messages\x1fmem\x1fsession-1") + mock_redis_client.delete.assert_called_once_with("13:chat_messages3:mem9:session-1", "chat_messages:session-1") async def test_clear_leaves_other_source_ids_untouched(self, mock_redis_client: MagicMock): with patch("agent_framework_redis._history_provider.redis.from_url") as mock_from_url: @@ -522,7 +563,9 @@ async def test_clear_leaves_other_source_ids_untouched(self, mock_redis_client: await audit.clear("session-1") # the destructive case from #7471: clearing one provider must not # delete the shared session's messages belonging to another provider - mock_redis_client.delete.assert_called_once_with("chat_messages\x1faudit\x1fsession-1") + mock_redis_client.delete.assert_called_once_with( + "13:chat_messages5:audit9:session-1", "chat_messages:session-1" + ) assert primary._redis_key("session-1") not in mock_redis_client.delete.call_args.args From 154cd8a1aadde7c68e3427e251fb5e2f6bb70e8b Mon Sep 17 00:00:00 2001 From: Yufeng He <40085740+he-yufeng@users.noreply.github.com> Date: Wed, 19 Aug 2026 04:41:22 +0800 Subject: [PATCH 4/4] fix: drop type ignores that pyright now rejects as unnecessary --- .../packages/redis/agent_framework_redis/_history_provider.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/python/packages/redis/agent_framework_redis/_history_provider.py b/python/packages/redis/agent_framework_redis/_history_provider.py index 4e24745550..002ba48624 100644 --- a/python/packages/redis/agent_framework_redis/_history_provider.py +++ b/python/packages/redis/agent_framework_redis/_history_provider.py @@ -153,9 +153,9 @@ async def get_messages( # atomic and no-ops if a concurrent write already landed there; the # legacy list then stays put and remains clearable via clear(). legacy_key = self._legacy_redis_key(session_id) - if legacy_key != key and await self._redis_client.exists(legacy_key): # type: ignore[misc] + if legacy_key != key and await self._redis_client.exists(legacy_key): with suppress(Exception): # a legacy key that vanished mid-read is a no-op - await self._redis_client.renamenx(legacy_key, key) # type: ignore[misc] + await self._redis_client.renamenx(legacy_key, key) redis_messages = await self._redis_client.lrange(key, 0, -1) # type: ignore[misc] messages: list[Message] = [] if redis_messages: