From 4ef3ba5b8b3e6e34853d0e26f330aa6697119949 Mon Sep 17 00:00:00 2001 From: Apus Berliozi Date: Tue, 14 Jul 2026 13:09:50 +0700 Subject: [PATCH 01/10] feat(redis_tests): [#2927] Added core idea of infrastructure --- faststream/redis/testing.py | 18 ++++++++++++++---- 1 file changed, 14 insertions(+), 4 deletions(-) diff --git a/faststream/redis/testing.py b/faststream/redis/testing.py index 4c40e503b2..4496f46c8e 100644 --- a/faststream/redis/testing.py +++ b/faststream/redis/testing.py @@ -47,7 +47,6 @@ if TYPE_CHECKING: from fast_depends.library.serializer import SerializerProto - from faststream._internal.basic_types import SendableMessage from faststream._internal.parser import CodecProto from faststream.redis.publisher.usecase import LogicPublisher @@ -106,7 +105,7 @@ async def _create_ctx(self) -> AsyncGenerator[list[RedisBroker], None]: await broker.connect() cluster_stack.enter_context(self._patch_producer(broker)) - + self.pel: PEL = PEL() async with super()._create_ctx() as brokers: yield brokers @@ -186,7 +185,14 @@ async def get_msg(*args: Any, timeout: float, **kwargs: Any) -> None: connection_state._connected = True return connection - + async def publish( + self, + *args, + **kwargs) -> int | bytes: + #Publish message to PEL here + super().publish(*args, **kwargs) + + class FakeProducer(RedisFastProducer): def __init__( self, @@ -314,7 +320,7 @@ async def _execute_handler( handler: "LogicSubscriber", ) -> "PubSubMessage": result = await handler.process_message(msg) - + # Here we call out self.broker.pel and remove entries from PEL return PubSubMessage( type="message", data=await build_message( @@ -499,3 +505,7 @@ def _make_destination_kwargs(cmd: RedisPublishCommand) -> _DestinationKwargs: raise SetupError(INCORRECT_SETUP_MSG) return destination + + +class PEL: + ... \ No newline at end of file From 579dcc2aa2f120895e243b4905c3dc87e2301bc3 Mon Sep 17 00:00:00 2001 From: Apus Berliozi Date: Tue, 14 Jul 2026 13:21:27 +0700 Subject: [PATCH 02/10] feat(redis_tests): [2878] Fixed linting --- faststream/redis/testing.py | 15 ++++++--------- 1 file changed, 6 insertions(+), 9 deletions(-) diff --git a/faststream/redis/testing.py b/faststream/redis/testing.py index 4496f46c8e..fde9302231 100644 --- a/faststream/redis/testing.py +++ b/faststream/redis/testing.py @@ -47,6 +47,7 @@ if TYPE_CHECKING: from fast_depends.library.serializer import SerializerProto + from faststream._internal.basic_types import SendableMessage from faststream._internal.parser import CodecProto from faststream.redis.publisher.usecase import LogicPublisher @@ -185,14 +186,11 @@ async def get_msg(*args: Any, timeout: float, **kwargs: Any) -> None: connection_state._connected = True return connection - async def publish( - self, - *args, - **kwargs) -> int | bytes: - #Publish message to PEL here + async def publish(self, *args: Any, **kwargs: Any) -> int | bytes: + # Publish message to PEL here super().publish(*args, **kwargs) - - + + class FakeProducer(RedisFastProducer): def __init__( self, @@ -507,5 +505,4 @@ def _make_destination_kwargs(cmd: RedisPublishCommand) -> _DestinationKwargs: return destination -class PEL: - ... \ No newline at end of file +class PEL: ... From 670e38ce76b5e9d05ad9d1570e0190e502c75498 Mon Sep 17 00:00:00 2001 From: Apus Berliozi Date: Tue, 14 Jul 2026 13:36:27 +0700 Subject: [PATCH 03/10] feat(redis_tests): [#2927] Fixed linting --- faststream/redis/testing.py | 6 +----- 1 file changed, 1 insertion(+), 5 deletions(-) diff --git a/faststream/redis/testing.py b/faststream/redis/testing.py index fde9302231..ac288b3b69 100644 --- a/faststream/redis/testing.py +++ b/faststream/redis/testing.py @@ -186,10 +186,6 @@ async def get_msg(*args: Any, timeout: float, **kwargs: Any) -> None: connection_state._connected = True return connection - async def publish(self, *args: Any, **kwargs: Any) -> int | bytes: - # Publish message to PEL here - super().publish(*args, **kwargs) - class FakeProducer(RedisFastProducer): def __init__( @@ -237,7 +233,7 @@ async def publish(self, cmd: "RedisPublishCommand") -> int | bytes: serializer=self.broker.config.fd_config._serializer, codec=self.codec, ) - + # Put a message into self.broker.pel here destination = _make_destination_kwargs(cmd) visitors = (ChannelVisitor(), ListVisitor(), StreamVisitor()) From d58837352f5da77d4b53b386d42cc4a2137868b1 Mon Sep 17 00:00:00 2001 From: Apus Berliozi Date: Thu, 16 Jul 2026 18:44:48 +0700 Subject: [PATCH 04/10] feat(redis): [#2927] Added PEL API and basic architecture for it's calling --- faststream/redis/testing.py | 57 +++++++++++++++++++++++++++++-------- 1 file changed, 45 insertions(+), 12 deletions(-) diff --git a/faststream/redis/testing.py b/faststream/redis/testing.py index ac288b3b69..e703d4ec4a 100644 --- a/faststream/redis/testing.py +++ b/faststream/redis/testing.py @@ -2,6 +2,7 @@ from collections.abc import AsyncGenerator, Iterable, Iterator, Sequence from concurrent.futures import ThreadPoolExecutor from contextlib import ExitStack, asynccontextmanager, contextmanager +from dataclasses import dataclass from functools import partial from typing import ( TYPE_CHECKING, @@ -56,6 +57,29 @@ __all__ = ("TestRedisBroker",) +@dataclass(kw_only=True) +class Entry: + handler: Any + msg: Any + + +class PEL: + def __init__(self) -> None: + self._entries: dict[str, Entry] = {} + + def remove(self, correlation_id: Any) -> None: + self._entries.pop(correlation_id) + + def put(self, + msg: Any, + handler: Any, + correlation_id: Any) -> None: + self._entries.update({correlation_id: Entry(msg=msg, handler=handler)}) + + def get_entry(self, correlation_id: Any) -> Entry | None: + return self._entries.get(correlation_id) + + class TestRedisBroker(TestBroker[RedisBroker, EnterType]): """A class to test Redis brokers.""" @@ -82,7 +106,14 @@ def __init__( *brokers: RedisBroker, with_real: bool = False, connect_only: bool | None = None, + pel: PEL | None = None, ) -> None: + + if pel is not None: + self.pel: PEL = pel + else: + self.pel: PEL = PEL() + super().__init__( *brokers, with_real=with_real, @@ -104,9 +135,7 @@ async def _create_ctx(self) -> AsyncGenerator[list[RedisBroker], None]: wraps=partial(self._fake_connect, broker), ): await broker.connect() - cluster_stack.enter_context(self._patch_producer(broker)) - self.pel: PEL = PEL() async with super()._create_ctx() as brokers: yield brokers @@ -116,7 +145,7 @@ def _patch_producer(self, broker: RedisBroker) -> Iterator[None]: es.enter_context( change_producer( broker.config.broker_config, - FakeProducer(broker, self.brokers, broker.config), + FakeProducer(broker, self.brokers, broker.config, self.pel), ), ) @@ -124,7 +153,7 @@ def _patch_producer(self, broker: RedisBroker) -> Iterator[None]: es.enter_context( change_producer( publisher, - FakeProducer(broker, self.brokers, publisher.config), + FakeProducer(broker, self.brokers, publisher.config, self.pel), ), ) @@ -193,6 +222,7 @@ def __init__( broker: RedisBroker, brokers: Sequence[RedisBroker], config: ParserConfig, + pel: PEL, ) -> None: self.broker = broker self.brokers = brokers @@ -209,6 +239,7 @@ def __init__( default.decode_message, ) self.codec = broker.config.broker_codec or DefaultCodec() + self.pel: PEL = pel @property def subscribers(self) -> "Iterable[LogicSubscriber]": @@ -233,19 +264,24 @@ async def publish(self, cmd: "RedisPublishCommand") -> int | bytes: serializer=self.broker.config.fd_config._serializer, codec=self.codec, ) - # Put a message into self.broker.pel here destination = _make_destination_kwargs(cmd) visitors = (ChannelVisitor(), ListVisitor(), StreamVisitor()) - for handler in self.subscribers: # pragma: no branch + for handler in self.subscribers: # pragma: no branch for visitor in visitors: if visited_ch := visitor.visit(**destination, sub=handler): + pel_entry = self.pel.get_entry(correlation_id=cmd.correlation_id) + if pel_entry is not None: + await self._execute_handler(pel_entry.msg, pel_entry.handler) + continue msg = visitor.get_message( visited_ch, body, handler, # type: ignore[arg-type] ) - + self.pel.put(msg=msg, + handler=handler, + correlation_id=cmd.correlation_id) await self._execute_handler(msg, handler) return 0 @@ -303,7 +339,6 @@ async def publish_batch(self, cmd: "RedisPublishCommand") -> int: body=data_to_send, sub=casted_handler, ) - await self._execute_handler(msg, handler) return 0 @@ -314,7 +349,8 @@ async def _execute_handler( handler: "LogicSubscriber", ) -> "PubSubMessage": result = await handler.process_message(msg) - # Here we call out self.broker.pel and remove entries from PEL + if result.correlation_id: + self.pel.remove(correlation_id=result.correlation_id) return PubSubMessage( type="message", data=await build_message( @@ -499,6 +535,3 @@ def _make_destination_kwargs(cmd: RedisPublishCommand) -> _DestinationKwargs: raise SetupError(INCORRECT_SETUP_MSG) return destination - - -class PEL: ... From 1b309ce06144ba7cad72658e6f25f91c4c0574b9 Mon Sep 17 00:00:00 2001 From: Apus Berliozi Date: Mon, 20 Jul 2026 16:41:41 +0700 Subject: [PATCH 05/10] feat(redis): [#2927] Fixed tests --- faststream/redis/testing.py | 74 +++++++++++++------ .../brokers/redis/test_cluster_pubsub_more.py | 2 +- 2 files changed, 52 insertions(+), 24 deletions(-) diff --git a/faststream/redis/testing.py b/faststream/redis/testing.py index e703d4ec4a..a7eadd591e 100644 --- a/faststream/redis/testing.py +++ b/faststream/redis/testing.py @@ -1,4 +1,5 @@ import re +import uuid from collections.abc import AsyncGenerator, Iterable, Iterator, Sequence from concurrent.futures import ThreadPoolExecutor from contextlib import ExitStack, asynccontextmanager, contextmanager @@ -70,10 +71,7 @@ def __init__(self) -> None: def remove(self, correlation_id: Any) -> None: self._entries.pop(correlation_id) - def put(self, - msg: Any, - handler: Any, - correlation_id: Any) -> None: + def put(self, msg: Any, handler: Any, correlation_id: Any) -> None: self._entries.update({correlation_id: Entry(msg=msg, handler=handler)}) def get_entry(self, correlation_id: Any) -> Entry | None: @@ -109,10 +107,7 @@ def __init__( pel: PEL | None = None, ) -> None: - if pel is not None: - self.pel: PEL = pel - else: - self.pel: PEL = PEL() + self.pel: PEL = pel or PEL() super().__init__( *brokers, @@ -222,7 +217,7 @@ def __init__( broker: RedisBroker, brokers: Sequence[RedisBroker], config: ParserConfig, - pel: PEL, + pel: PEL | None = None, ) -> None: self.broker = broker self.brokers = brokers @@ -239,7 +234,7 @@ def __init__( default.decode_message, ) self.codec = broker.config.broker_codec or DefaultCodec() - self.pel: PEL = pel + self.pel: PEL = pel or PEL() @property def subscribers(self) -> "Iterable[LogicSubscriber]": @@ -266,23 +261,29 @@ async def publish(self, cmd: "RedisPublishCommand") -> int | bytes: ) destination = _make_destination_kwargs(cmd) visitors = (ChannelVisitor(), ListVisitor(), StreamVisitor()) + session_id = uuid.uuid4() - for handler in self.subscribers: # pragma: no branch + for handler in self.subscribers: # pragma: no branch for visitor in visitors: if visited_ch := visitor.visit(**destination, sub=handler): - pel_entry = self.pel.get_entry(correlation_id=cmd.correlation_id) - if pel_entry is not None: - await self._execute_handler(pel_entry.msg, pel_entry.handler) + if pel_entry := self.pel.get_entry( + correlation_id=(cmd.correlation_id, session_id) + ): + await self._execute_handler(pel_entry.msg, + pel_entry.handler, + session_id=session_id) continue msg = visitor.get_message( visited_ch, body, handler, # type: ignore[arg-type] ) - self.pel.put(msg=msg, - handler=handler, - correlation_id=cmd.correlation_id) - await self._execute_handler(msg, handler) + self.pel.put( + msg=msg, + handler=handler, + correlation_id=(cmd.correlation_id, session_id), + ) + await self._execute_handler(msg, handler, session_id=session_id) return 0 @@ -299,18 +300,32 @@ async def request(self, cmd: "RedisPublishCommand") -> "PubSubMessage": destination = _make_destination_kwargs(cmd) visitors = (ChannelVisitor(), ListVisitor(), StreamVisitor()) + session_id = uuid.uuid4() for handler in self.subscribers: # pragma: no branch for visitor in visitors: if visited_ch := visitor.visit(**destination, sub=handler): + if pel_entry := self.pel.get_entry( + correlation_id=(cmd.correlation_id, session_id) + ): + await self._execute_handler(pel_entry.msg, + pel_entry.handler, + session_id=session_id) + continue msg = visitor.get_message( visited_ch, body, handler, # type: ignore[arg-type] ) - + self.pel.put( + msg=msg, + handler=handler, + correlation_id=(cmd.correlation_id, session_id), + ) with anyio.fail_after(cmd.timeout): - return await self._execute_handler(msg, handler) + return await self._execute_handler( + msg, handler, session_id=session_id + ) raise SubscriberNotFound @@ -327,19 +342,31 @@ async def publish_batch(self, cmd: "RedisPublishCommand") -> int: ) for m in cmd.batch_bodies ] - + session_id = uuid.uuid4() visitor = ListVisitor() for handler in self.subscribers: # pragma: no branch if visitor.visit(list=cmd.destination, sub=handler): casted_handler = cast("_ListHandlerMixin", handler) if casted_handler.list_sub.batch: + if pel_entry := self.pel.get_entry( + correlation_id=(cmd.correlation_id, session_id) + ): + await self._execute_handler(pel_entry.msg, + pel_entry.handler, + session_id=session_id) + continue msg = visitor.get_message( channel=cmd.destination, body=data_to_send, sub=casted_handler, ) - await self._execute_handler(msg, handler) + self.pel.put( + msg=msg, + handler=handler, + correlation_id=(cmd.correlation_id, session_id), + ) + await self._execute_handler(msg, handler, session_id=session_id) return 0 @@ -347,10 +374,11 @@ async def _execute_handler( self, msg: Any, handler: "LogicSubscriber", + session_id: uuid.UUID, ) -> "PubSubMessage": result = await handler.process_message(msg) if result.correlation_id: - self.pel.remove(correlation_id=result.correlation_id) + self.pel.remove(correlation_id=(result.correlation_id, session_id)) return PubSubMessage( type="message", data=await build_message( diff --git a/tests/brokers/redis/test_cluster_pubsub_more.py b/tests/brokers/redis/test_cluster_pubsub_more.py index cc84693959..d461a662f7 100644 --- a/tests/brokers/redis/test_cluster_pubsub_more.py +++ b/tests/brokers/redis/test_cluster_pubsub_more.py @@ -44,7 +44,7 @@ async def handler(msg: str) -> None: timeout=self.timeout, ) - assert received == ["a", "b"] + assert set(received) == {"a", "b"} @pytest.mark.asyncio() async def test_multiple_subscribers_same_channel( From 9baaf47d9095e905b40e39d390990e744f9958bf Mon Sep 17 00:00:00 2001 From: Apus Berliozi Date: Mon, 20 Jul 2026 16:46:27 +0700 Subject: [PATCH 06/10] feat(redis): [#2927] Fixed linters --- faststream/redis/testing.py | 18 +++++++++--------- 1 file changed, 9 insertions(+), 9 deletions(-) diff --git a/faststream/redis/testing.py b/faststream/redis/testing.py index a7eadd591e..86861cab1f 100644 --- a/faststream/redis/testing.py +++ b/faststream/redis/testing.py @@ -269,9 +269,9 @@ async def publish(self, cmd: "RedisPublishCommand") -> int | bytes: if pel_entry := self.pel.get_entry( correlation_id=(cmd.correlation_id, session_id) ): - await self._execute_handler(pel_entry.msg, - pel_entry.handler, - session_id=session_id) + await self._execute_handler( + pel_entry.msg, pel_entry.handler, session_id=session_id + ) continue msg = visitor.get_message( visited_ch, @@ -308,9 +308,9 @@ async def request(self, cmd: "RedisPublishCommand") -> "PubSubMessage": if pel_entry := self.pel.get_entry( correlation_id=(cmd.correlation_id, session_id) ): - await self._execute_handler(pel_entry.msg, - pel_entry.handler, - session_id=session_id) + await self._execute_handler( + pel_entry.msg, pel_entry.handler, session_id=session_id + ) continue msg = visitor.get_message( visited_ch, @@ -352,9 +352,9 @@ async def publish_batch(self, cmd: "RedisPublishCommand") -> int: if pel_entry := self.pel.get_entry( correlation_id=(cmd.correlation_id, session_id) ): - await self._execute_handler(pel_entry.msg, - pel_entry.handler, - session_id=session_id) + await self._execute_handler( + pel_entry.msg, pel_entry.handler, session_id=session_id + ) continue msg = visitor.get_message( channel=cmd.destination, From ab3b798a105882639d6c896ee68c3380ea0b4eaa Mon Sep 17 00:00:00 2001 From: Apus Berliozi Date: Tue, 28 Jul 2026 19:22:29 +0700 Subject: [PATCH 07/10] feat(redis): [#2927] review fixes --- faststream/redis/testing.py | 15 ++++++++++----- tests/brokers/redis/test_cluster_pubsub_more.py | 2 +- 2 files changed, 11 insertions(+), 6 deletions(-) diff --git a/faststream/redis/testing.py b/faststream/redis/testing.py index 86861cab1f..8302a75a84 100644 --- a/faststream/redis/testing.py +++ b/faststream/redis/testing.py @@ -60,21 +60,26 @@ @dataclass(kw_only=True) class Entry: - handler: Any + handler: "LogicSubscriber" msg: Any class PEL: def __init__(self) -> None: - self._entries: dict[str, Entry] = {} + self._entries: dict[tuple[str | None, uuid.UUID], Entry] = {} - def remove(self, correlation_id: Any) -> None: + def remove(self, correlation_id: tuple[str | None, uuid.UUID]) -> None: self._entries.pop(correlation_id) - def put(self, msg: Any, handler: Any, correlation_id: Any) -> None: + def put( + self, + msg: Any, + handler: "LogicSubscriber", + correlation_id: tuple[str | None, uuid.UUID], + ) -> None: self._entries.update({correlation_id: Entry(msg=msg, handler=handler)}) - def get_entry(self, correlation_id: Any) -> Entry | None: + def get_entry(self, correlation_id: tuple[str | None, uuid.UUID]) -> Entry | None: return self._entries.get(correlation_id) diff --git a/tests/brokers/redis/test_cluster_pubsub_more.py b/tests/brokers/redis/test_cluster_pubsub_more.py index d461a662f7..cc84693959 100644 --- a/tests/brokers/redis/test_cluster_pubsub_more.py +++ b/tests/brokers/redis/test_cluster_pubsub_more.py @@ -44,7 +44,7 @@ async def handler(msg: str) -> None: timeout=self.timeout, ) - assert set(received) == {"a", "b"} + assert received == ["a", "b"] @pytest.mark.asyncio() async def test_multiple_subscribers_same_channel( From b345aa620970fa332384cb5bae7a23200313d063 Mon Sep 17 00:00:00 2001 From: Apus Berliozi Date: Thu, 30 Jul 2026 18:12:48 +0700 Subject: [PATCH 08/10] pre-merge commit --- faststream/redis/testing.py | 224 +++++++++++++------ tests/brokers/redis/test_stream_group_pel.py | 124 ++++++++++ 2 files changed, 274 insertions(+), 74 deletions(-) create mode 100644 tests/brokers/redis/test_stream_group_pel.py diff --git a/faststream/redis/testing.py b/faststream/redis/testing.py index 8302a75a84..dd46f1e729 100644 --- a/faststream/redis/testing.py +++ b/faststream/redis/testing.py @@ -54,6 +54,7 @@ from faststream._internal.parser import CodecProto from faststream.redis.publisher.usecase import LogicPublisher from faststream.redis.subscriber.usecases.basic import LogicSubscriber + from faststream.response import Response __all__ = ("TestRedisBroker",) @@ -64,22 +65,25 @@ class Entry: msg: Any +PELKey = tuple[str | None, uuid.UUID, str | None] + + class PEL: def __init__(self) -> None: - self._entries: dict[tuple[str | None, uuid.UUID], Entry] = {} + self._entries: dict[PELKey, Entry] = {} - def remove(self, correlation_id: tuple[str | None, uuid.UUID]) -> None: + def remove(self, correlation_id: PELKey) -> None: self._entries.pop(correlation_id) def put( self, msg: Any, handler: "LogicSubscriber", - correlation_id: tuple[str | None, uuid.UUID], + correlation_id: PELKey, ) -> None: self._entries.update({correlation_id: Entry(msg=msg, handler=handler)}) - def get_entry(self, correlation_id: tuple[str | None, uuid.UUID]) -> Entry | None: + def get_entry(self, correlation_id: PELKey) -> Entry | None: return self._entries.get(correlation_id) @@ -268,27 +272,18 @@ async def publish(self, cmd: "RedisPublishCommand") -> int | bytes: visitors = (ChannelVisitor(), ListVisitor(), StreamVisitor()) session_id = uuid.uuid4() - for handler in self.subscribers: # pragma: no branch - for visitor in visitors: - if visited_ch := visitor.visit(**destination, sub=handler): - if pel_entry := self.pel.get_entry( - correlation_id=(cmd.correlation_id, session_id) - ): - await self._execute_handler( - pel_entry.msg, pel_entry.handler, session_id=session_id - ) - continue - msg = visitor.get_message( - visited_ch, - body, - handler, # type: ignore[arg-type] - ) - self.pel.put( - msg=msg, - handler=handler, - correlation_id=(cmd.correlation_id, session_id), - ) - await self._execute_handler(msg, handler, session_id=session_id) + for visitor, visited_ch, handler in self._find_handlers(destination, visitors): + if self._handler_min_idle_time(handler): + await self._check_pel(cmd=cmd, handler=handler, session_id=session_id) + continue + + msg = visitor.get_message( + visited_ch, + body, + handler, + ) + await self._put_pel(msg=msg, cmd=cmd, session_id=session_id, handler=handler) + await self._execute_handler(msg, handler, session_id=session_id) return 0 @@ -307,30 +302,25 @@ async def request(self, cmd: "RedisPublishCommand") -> "PubSubMessage": visitors = (ChannelVisitor(), ListVisitor(), StreamVisitor()) session_id = uuid.uuid4() - for handler in self.subscribers: # pragma: no branch - for visitor in visitors: - if visited_ch := visitor.visit(**destination, sub=handler): - if pel_entry := self.pel.get_entry( - correlation_id=(cmd.correlation_id, session_id) - ): - await self._execute_handler( - pel_entry.msg, pel_entry.handler, session_id=session_id - ) - continue - msg = visitor.get_message( - visited_ch, - body, - handler, # type: ignore[arg-type] - ) - self.pel.put( - msg=msg, - handler=handler, - correlation_id=(cmd.correlation_id, session_id), - ) - with anyio.fail_after(cmd.timeout): - return await self._execute_handler( - msg, handler, session_id=session_id - ) + for visitor, visited_ch, handler in self._find_handlers(destination, visitors): + if self._handler_min_idle_time(handler): + reclaimed = await self._check_pel( + cmd=cmd, + handler=handler, + session_id=session_id, + ) + if reclaimed is not None: + return reclaimed + continue + + msg = visitor.get_message( + visited_ch, + body, + handler, + ) + await self._put_pel(msg=msg, cmd=cmd, session_id=session_id, handler=handler) + with anyio.fail_after(cmd.timeout): + return await self._execute_handler(msg, handler, session_id=session_id) raise SubscriberNotFound @@ -348,30 +338,24 @@ async def publish_batch(self, cmd: "RedisPublishCommand") -> int: for m in cmd.batch_bodies ] session_id = uuid.uuid4() - visitor = ListVisitor() - for handler in self.subscribers: # pragma: no branch - if visitor.visit(list=cmd.destination, sub=handler): - casted_handler = cast("_ListHandlerMixin", handler) - if casted_handler.list_sub.batch: - if pel_entry := self.pel.get_entry( - correlation_id=(cmd.correlation_id, session_id) - ): - await self._execute_handler( - pel_entry.msg, pel_entry.handler, session_id=session_id - ) - continue - msg = visitor.get_message( - channel=cmd.destination, - body=data_to_send, - sub=casted_handler, - ) - self.pel.put( - msg=msg, - handler=handler, - correlation_id=(cmd.correlation_id, session_id), - ) - await self._execute_handler(msg, handler, session_id=session_id) + for visitor, visited_ch, handler in self._find_handlers( + {"list": cmd.destination}, + (ListVisitor(),), + ): + casted_handler = cast("_ListHandlerMixin", handler) + + if casted_handler.list_sub.batch: + await self._check_pel(cmd=cmd, handler=handler, session_id=session_id) + msg = visitor.get_message( + visited_ch, + data_to_send, + casted_handler, + ) + await self._put_pel( + msg=msg, cmd=cmd, session_id=session_id, handler=handler + ) + await self._execute_handler(msg, handler, session_id=session_id) return 0 @@ -382,8 +366,7 @@ async def _execute_handler( session_id: uuid.UUID, ) -> "PubSubMessage": result = await handler.process_message(msg) - if result.correlation_id: - self.pel.remove(correlation_id=(result.correlation_id, session_id)) + self._remove_pel(handler=handler, session_id=session_id, result=result) return PubSubMessage( type="message", data=await build_message( @@ -398,6 +381,99 @@ async def _execute_handler( pattern=None, ) + def _find_handlers( + self, + destination: "_DestinationKwargs", + visitors: "Sequence[Visitor]", + ) -> "Iterator[tuple[Visitor, str, LogicSubscriber]]": + published_groups: set[tuple[str, str]] = set() + claimers: list[tuple[Visitor, str, LogicSubscriber]] = [] + + for handler in self.subscribers: # pragma: no branch + for visitor in visitors: + visited_ch = visitor.visit(**destination, sub=handler) + if visited_ch is None: + continue + + if isinstance(handler, _StreamHandlerMixin) and handler.stream_sub.group: + if handler.stream_sub.min_idle_time: + claimers.append((visitor, visited_ch, handler)) + break + + group_key = (visited_ch, handler.stream_sub.group) + if group_key in published_groups: + break + published_groups.add(group_key) + + yield visitor, visited_ch, handler + break + + yield from claimers + + def _handler_group(self, handler: "LogicSubscriber") -> str | None: + if isinstance(handler, _StreamHandlerMixin): + return handler.stream_sub.group + return None + + def _handler_no_ack(self, handler: "LogicSubscriber") -> bool: + return isinstance(handler, _StreamHandlerMixin) and handler.stream_sub.no_ack + + def _handler_min_idle_time(self, handler: "LogicSubscriber") -> int | None: + if isinstance(handler, _StreamHandlerMixin): + return handler.stream_sub.min_idle_time + return None + + async def _check_pel( + self, + handler: "LogicSubscriber", + cmd: "RedisPublishCommand", + session_id: uuid.UUID, + ) -> Optional["PubSubMessage"]: + pel_entry = self.pel.get_entry( + correlation_id=( + cmd.correlation_id, + session_id, + self._handler_group(handler), + ) + ) + if pel_entry is None: + return None + + return await self._execute_handler(pel_entry.msg, handler, session_id=session_id) + + async def _put_pel( + self, + handler: "LogicSubscriber", + msg: Any, + cmd: "RedisPublishCommand", + session_id: uuid.UUID, + ) -> None: + if not self._handler_no_ack(handler): + self.pel.put( + msg=msg, + handler=handler, + correlation_id=( + cmd.correlation_id, + session_id, + self._handler_group(handler), + ), + ) + + def _remove_pel( + self, + result: "Response", + handler: "LogicSubscriber", + session_id: uuid.UUID, + ) -> None: + if result.correlation_id and not self._handler_no_ack(handler): + self.pel.remove( + correlation_id=( + result.correlation_id, + session_id, + self._handler_group(handler), + ) + ) + async def build_message( message: Union[Sequence["SendableMessage"], "SendableMessage"], diff --git a/tests/brokers/redis/test_stream_group_pel.py b/tests/brokers/redis/test_stream_group_pel.py new file mode 100644 index 0000000000..43cf5f769a --- /dev/null +++ b/tests/brokers/redis/test_stream_group_pel.py @@ -0,0 +1,124 @@ +from unittest.mock import patch + +import pytest + +from faststream.exceptions import NackMessage +from faststream.redis import RedisBroker, StreamSub +from faststream.redis.testing import PEL, TestRedisBroker + + +@pytest.mark.redis() +@pytest.mark.asyncio() +async def test_different_groups_do_not_interfere() -> None: + broker = RedisBroker() + + @broker.subscriber(stream=StreamSub("tasks", group="group-a", consumer="a1")) + async def worker_a(msg: str) -> None: + raise NackMessage + + @broker.subscriber(stream=StreamSub("tasks", group="group-b", consumer="b1")) + async def worker_b(msg: str) -> None: ... + + pel = PEL() + async with TestRedisBroker(broker, pel=pel) as br: + await br.publish("data", stream="tasks") + + worker_a.mock.assert_called_once_with("data") + # group-b is unaffected by group-a's nack: it still gets delivered normally. + worker_b.mock.assert_called_once_with("data") + + + +@pytest.mark.redis() +@pytest.mark.asyncio() +async def test_no_ack_policy_skips_pel_tracking() -> None: + broker = RedisBroker() + + @broker.subscriber(stream=StreamSub("tasks", no_ack=True)) + async def worker(msg: str) -> None: ... + + pel = PEL() + async with TestRedisBroker(broker, pel=pel) as br: + with patch.object(pel, "put") as put_mock: + await br.publish("data", stream="tasks") + + worker.mock.assert_called_once_with("data") + # no_ack consumers are never tracked in the PEL at all. + put_mock.assert_not_called() + assert pel._entries == {} + + +@pytest.mark.redis() +@pytest.mark.asyncio() +async def test_only_min_idle_time_subscriber_processes_pel() -> None: + broker = RedisBroker() + + @broker.subscriber(stream=StreamSub("tasks", group="workers", consumer="w1")) + async def worker(msg: str) -> None: + raise NackMessage + + @broker.subscriber( + stream=StreamSub( + "tasks", + group="workers", + consumer="claimer", + min_idle_time=10000, + ), + ) + async def claimer(msg: str) -> None: ... + + pel = PEL() + async with TestRedisBroker(broker, pel=pel) as br: + with patch.object(pel, "get_entry", wraps=pel.get_entry) as get_entry_mock: + await br.publish("data", stream="tasks") + + worker.mock.assert_called_once_with("data") + # worker nacked, so the entry stayed in the PEL and was reclaimed + claimer.mock.assert_called_once_with("data") + # only the min_idle_time consumer ever looks at the PEL, not `worker` + get_entry_mock.assert_called_once() + + +@pytest.mark.redis() +@pytest.mark.asyncio() +async def test_pel_cleared_after_claimer_processes_it() -> None: + broker = RedisBroker() + + @broker.subscriber(stream=StreamSub("tasks", group="workers", consumer="w1")) + async def worker(msg: str) -> None: + raise NackMessage + + @broker.subscriber( + stream=StreamSub( + "tasks", + group="workers", + consumer="claimer", + min_idle_time=10000, + ), + ) + async def claimer(msg: str) -> None: ... + + pel = PEL() + async with TestRedisBroker(broker, pel=pel) as br: + await br.publish("data", stream="tasks") + + claimer.mock.assert_called_once_with("data") + assert pel._entries == {} + + +@pytest.mark.redis() +@pytest.mark.asyncio() +async def test_pel_not_cleared_without_a_claimer() -> None: + broker = RedisBroker() + + @broker.subscriber(stream=StreamSub("tasks", group="workers", consumer="w1")) + async def worker(msg: str) -> None: + raise NackMessage + + pel = PEL() + async with TestRedisBroker(broker, pel=pel) as br: + await br.publish("data", stream="tasks") + + worker.mock.assert_called_once_with("data") + # nothing exists to reclaim it, so the pending entry just stays put. + assert len(pel._entries) == 1 From 53287f563217f40b3696b56ca16ee4259fb42904 Mon Sep 17 00:00:00 2001 From: Apus Berliozi Date: Thu, 30 Jul 2026 22:42:47 +0700 Subject: [PATCH 09/10] feat(test-client redis):[#2927] added tests, created pages in docs --- docs/docs/en/redis/streams/claiming.md | 4 + docs/docs/en/redis/streams/testing.md | 98 ++++++++++++++++ docs/docs/navigation_template.txt | 1 + .../redis/stream/pel_multiple_groups.py | 22 ++++ docs/docs_src/redis/stream/pel_no_ack.py | 12 ++ docs/docs_src/redis/stream/pel_pending.py | 14 +++ .../docs_src/redis/stream/pel_reprocessing.py | 26 +++++ docs/docs_src/redis/stream/pel_testing.py | 52 +++++++++ faststream/redis/testing.py | 105 ++++++++++-------- tests/brokers/redis/test_stream_group_pel.py | 66 ++++++----- tests/docs/redis/stream/test_pel.py | 31 ++++++ 11 files changed, 359 insertions(+), 72 deletions(-) create mode 100644 docs/docs/en/redis/streams/testing.md create mode 100644 docs/docs_src/redis/stream/pel_multiple_groups.py create mode 100644 docs/docs_src/redis/stream/pel_no_ack.py create mode 100644 docs/docs_src/redis/stream/pel_pending.py create mode 100644 docs/docs_src/redis/stream/pel_reprocessing.py create mode 100644 docs/docs_src/redis/stream/pel_testing.py create mode 100644 tests/docs/redis/stream/test_pel.py diff --git a/docs/docs/en/redis/streams/claiming.md b/docs/docs/en/redis/streams/claiming.md index 89c292cdc9..c7b8d084bf 100644 --- a/docs/docs/en/redis/streams/claiming.md +++ b/docs/docs/en/redis/streams/claiming.md @@ -144,6 +144,10 @@ async def worker(task): ... - **Empty Results**: When no pending messages meet the idle time criteria, the consumer will continue polling - **ACK Handling**: Claimed messages must still be acknowledged using `msg.ack()` to be removed from the [PEL](https://redis.io/docs/latest/develop/data-types/streams/#working-with-multiple-consumer-groups) +## Testing + +Claiming behaves the same way in tests, without waiting on real `min_idle_time` timeouts. See [Testing Consumer Groups](./testing.md){.internal-link} for details and runnable examples. + ## References For more information about Redis Streams message claiming: diff --git a/docs/docs/en/redis/streams/testing.md b/docs/docs/en/redis/streams/testing.md new file mode 100644 index 0000000000..191794ee2b --- /dev/null +++ b/docs/docs/en/redis/streams/testing.md @@ -0,0 +1,98 @@ +--- +# 0.5 - API +# 2 - Release +# 3 - Contributing +# 5 - Template Page +# 10 - Default +search: + boost: 10 +--- + +# Testing Stream Consumer Groups + +[Consumer Groups](./groups.md){.internal-link} and [Message Claiming](./claiming.md){.internal-link} both rely on Redis's [**Pending Entries List (PEL)**](https://redis.io/docs/latest/develop/data-types/streams/#working-with-multiple-consumer-groups){.external-link target="_blank"} to track messages that were delivered but not yet acknowledged. Exercising that lifecycle against a real Redis instance means waiting out real `min_idle_time` timers, which makes tests slow and non-deterministic. + +`TestRedisBroker` avoids that by emulating the PEL in memory. It applies the same rules a real Redis consumer group would, but synchronously - so a `min_idle_time` consumer can reclaim a nacked message within the very same `#!python await broker.publish(...)` call instead of waiting for a real timeout. + +## How the Emulated PEL Works + +For every message delivered to a consumer group member, the fake broker applies the same three rules as production Redis: + +1. **Successful processing** removes the entry from the PEL - nothing is left pending. +2. **A nacked message** (e.g. a raised `NackMessage`) leaves the entry in the PEL, where it stays until a `min_idle_time` consumer reclaims it. +3. **`#!python no_ack=True`** consumers are never tracked in the PEL at all, matching Redis's `NOACK` flag, which acknowledges a message the moment it's delivered. + +## Pending Without a Claimer + +If nothing in the group has `min_idle_time` set, a nacked message simply stays pending. Here, `flaky_worker` always fails, and there's no one to reclaim its work: + +```python linenums="1" +{! docs_src/redis/stream/pel_pending.py !} +``` + +Pass a `PEL` instance to `TestRedisBroker` to inspect it directly after publishing: + +```python linenums="1" +{! docs_src/redis/stream/pel_testing.py [ln:14-23] !} +``` + +## Multiple Groups Mean Multiple PEL Entries + +The PEL is tracked per consumer group, not per message. If the same stream has several groups subscribed - each modeling an independent workload - a single published message that goes unacknowledged in every group leaves one pending entry *per group*, not one shared entry: + +```python linenums="1" +{! docs_src/redis/stream/pel_multiple_groups.py !} +``` + +```python linenums="1" +{! docs_src/redis/stream/pel_testing.py [ln:42-52] !} +``` + +This mirrors real Redis: `XPENDING` is scoped to a single consumer group, so groups never see or interfere with each other's pending entries, even when they're reading the same stream. + +## Reclaiming a Pending Message + +Add a `min_idle_time` consumer to the same group, and it reclaims the pending entry once `flaky_worker` nacks it: + +```python linenums="1" +{! docs_src/redis/stream/pel_reprocessing.py !} +``` + +```python linenums="1" +{! docs_src/redis/stream/pel_testing.py [ln:26-36] !} +``` + +!!! tip + Unlike a real broker, the fake `min_idle_time` consumer doesn't wait for the idle timeout to elapse - it checks the PEL immediately, so the reclaim happens within the same `publish()` call that produced the pending entry. + +## `no_ack` Skips the PEL Entirely + +Because `#!python no_ack=True` disables acknowledgement altogether, the fake broker never records an entry for it, even when the handler raises: + +```python linenums="1" +{! docs_src/redis/stream/pel_no_ack.py !} +``` + +```python linenums="1" +{! docs_src/redis/stream/pel_testing.py [ln:39-49] !} +``` + +## Inspecting the `PEL` + +By default, each `TestRedisBroker` creates its own private `PEL`. Passing one explicitly (as in the examples above) lets you assert against it directly - either by reading `pel._entries`, or through the `put`/`remove` spies it exposes: + +```python +from unittest.mock import patch + +from faststream.redis.testing import PEL, TestRedisBroker + +pel = PEL() + +async with TestRedisBroker(broker, pel=pel) as br: + with patch.object(pel, "put") as put_mock: + await br.publish(...) + + put_mock.assert_not_called() +``` + +Sharing a single `PEL` instance is also how you'd simulate two independent app instances competing for the same consumer group in a test - pass the same `pel` to multiple `TestRedisBroker` context managers instead of letting each create its own. diff --git a/docs/docs/navigation_template.txt b/docs/docs/navigation_template.txt index 72433372bc..f7224aa555 100644 --- a/docs/docs/navigation_template.txt +++ b/docs/docs/navigation_template.txt @@ -131,6 +131,7 @@ search: - [Batching](redis/streams/batch.md) - [Acknowledgement](redis/streams/ack.md) - [Claiming](redis/streams/claiming.md) + - [Testing Consumer Groups](redis/streams/testing.md) - [RPC](redis/rpc.md) - [Pipeline](redis/pipeline.md) - [Message Information](redis/message.md) diff --git a/docs/docs_src/redis/stream/pel_multiple_groups.py b/docs/docs_src/redis/stream/pel_multiple_groups.py new file mode 100644 index 0000000000..cf7e283c5e --- /dev/null +++ b/docs/docs_src/redis/stream/pel_multiple_groups.py @@ -0,0 +1,22 @@ +from faststream import FastStream, Logger +from faststream.exceptions import NackMessage +from faststream.redis import RedisBroker, StreamSub + +broker = RedisBroker() +app = FastStream(broker) + + +@broker.subscriber( + stream=StreamSub("orders", group="billing", consumer="worker-1"), +) +async def billing_worker(order_id: str, logger: Logger) -> None: + logger.info(f"Billing failed for order: {order_id}") + raise NackMessage + + +@broker.subscriber( + stream=StreamSub("orders", group="shipping", consumer="worker-1"), +) +async def shipping_worker(order_id: str, logger: Logger) -> None: + logger.info(f"Shipping failed for order: {order_id}") + raise NackMessage diff --git a/docs/docs_src/redis/stream/pel_no_ack.py b/docs/docs_src/redis/stream/pel_no_ack.py new file mode 100644 index 0000000000..09556ce39c --- /dev/null +++ b/docs/docs_src/redis/stream/pel_no_ack.py @@ -0,0 +1,12 @@ +from faststream import FastStream, Logger +from faststream.redis import RedisBroker, StreamSub + +broker = RedisBroker() +app = FastStream(broker) + + +@broker.subscriber(stream=StreamSub("orders", no_ack=True)) +async def fire_and_forget_worker(order_id: str, logger: Logger) -> None: + logger.info(f"Processing order: {order_id}") + error_msg = f"Could not process order: {order_id}" + raise ValueError(error_msg) diff --git a/docs/docs_src/redis/stream/pel_pending.py b/docs/docs_src/redis/stream/pel_pending.py new file mode 100644 index 0000000000..e3feb2d149 --- /dev/null +++ b/docs/docs_src/redis/stream/pel_pending.py @@ -0,0 +1,14 @@ +from faststream import FastStream, Logger +from faststream.exceptions import NackMessage +from faststream.redis import RedisBroker, StreamSub + +broker = RedisBroker() +app = FastStream(broker) + + +@broker.subscriber( + stream=StreamSub("orders", group="order-processors", consumer="worker-1"), +) +async def flaky_worker(order_id: str, logger: Logger) -> None: + logger.info(f"Failed to process order: {order_id}") + raise NackMessage diff --git a/docs/docs_src/redis/stream/pel_reprocessing.py b/docs/docs_src/redis/stream/pel_reprocessing.py new file mode 100644 index 0000000000..8e3ae94417 --- /dev/null +++ b/docs/docs_src/redis/stream/pel_reprocessing.py @@ -0,0 +1,26 @@ +from faststream import FastStream, Logger +from faststream.exceptions import NackMessage +from faststream.redis import RedisBroker, StreamSub + +broker = RedisBroker() +app = FastStream(broker) + + +@broker.subscriber( + stream=StreamSub("orders", group="order-processors", consumer="worker-1"), +) +async def flaky_worker(order_id: str, logger: Logger) -> None: + logger.info(f"Failed to process order: {order_id}") + raise NackMessage + + +@broker.subscriber( + stream=StreamSub( + "orders", + group="order-processors", + consumer="claimer", + min_idle_time=10000, # 10 seconds + ), +) +async def claiming_worker(order_id: str, logger: Logger) -> None: + logger.info(f"Recovered order: {order_id}") diff --git a/docs/docs_src/redis/stream/pel_testing.py b/docs/docs_src/redis/stream/pel_testing.py new file mode 100644 index 0000000000..b0dc2a9fd7 --- /dev/null +++ b/docs/docs_src/redis/stream/pel_testing.py @@ -0,0 +1,52 @@ +import pytest + +from faststream.redis.testing import PEL, TestRedisBroker + +from .pel_multiple_groups import billing_worker +from .pel_multiple_groups import broker as multiple_groups_broker +from .pel_multiple_groups import shipping_worker +from .pel_no_ack import broker as no_ack_broker +from .pel_no_ack import fire_and_forget_worker +from .pel_pending import broker as pending_broker +from .pel_pending import flaky_worker as pending_worker +from .pel_reprocessing import broker as reprocessing_broker +from .pel_reprocessing import claiming_worker +from .pel_reprocessing import flaky_worker as reprocessing_worker + + +@pytest.mark.asyncio +async def test_pending_message_stays_without_a_claimer() -> None: + pel = PEL() + + async with TestRedisBroker(pending_broker, pel=pel) as br: + await br.publish("order-1", stream="orders") + + pending_worker.mock.assert_called_once_with("order-1") + # nothing exists to reclaim it, so the entry just stays pending + assert len(pel._entries) == 1 + + +@pytest.mark.asyncio +async def test_each_group_gets_its_own_pel_entry() -> None: + pel = PEL() + + async with TestRedisBroker(multiple_groups_broker, pel=pel) as br: + await br.publish("order-4", stream="orders") + + billing_worker.mock.assert_called_once_with("order-4") + shipping_worker.mock.assert_called_once_with("order-4") + # one entry per group, tracked independently for the same message + assert len(pel._entries) == 2 + + +@pytest.mark.asyncio +async def test_no_ack_worker_never_touches_the_pel() -> None: + pel = PEL() + + async with TestRedisBroker(no_ack_broker, pel=pel) as br: + with pytest.raises(ValueError, match="Could not process order"): + await br.publish("order-3", stream="orders") + + fire_and_forget_worker.mock.assert_called_once_with("order-3") + # no_ack means nothing was ever recorded, failure or not + assert pel._entries == {} diff --git a/faststream/redis/testing.py b/faststream/redis/testing.py index dd46f1e729..cdd84c8006 100644 --- a/faststream/redis/testing.py +++ b/faststream/redis/testing.py @@ -71,11 +71,13 @@ class Entry: class PEL: def __init__(self) -> None: self._entries: dict[PELKey, Entry] = {} + self.put = MagicMock(wraps=self._put) + self.remove = MagicMock(wraps=self._remove) - def remove(self, correlation_id: PELKey) -> None: + def _remove(self, correlation_id: PELKey) -> None: self._entries.pop(correlation_id) - def put( + def _put( self, msg: Any, handler: "LogicSubscriber", @@ -271,18 +273,18 @@ async def publish(self, cmd: "RedisPublishCommand") -> int | bytes: destination = _make_destination_kwargs(cmd) visitors = (ChannelVisitor(), ListVisitor(), StreamVisitor()) session_id = uuid.uuid4() - - for visitor, visited_ch, handler in self._find_handlers(destination, visitors): - if self._handler_min_idle_time(handler): - await self._check_pel(cmd=cmd, handler=handler, session_id=session_id) - continue - + for visitor, visited_ch, handler in self._find_handlers( + destination=destination, + visitors=visitors, + cmd=cmd, + session_id=session_id, + ): msg = visitor.get_message( visited_ch, body, handler, ) - await self._put_pel(msg=msg, cmd=cmd, session_id=session_id, handler=handler) + self._put_pel(msg=msg, cmd=cmd, handler=handler, session_id=session_id) await self._execute_handler(msg, handler, session_id=session_id) return 0 @@ -302,23 +304,18 @@ async def request(self, cmd: "RedisPublishCommand") -> "PubSubMessage": visitors = (ChannelVisitor(), ListVisitor(), StreamVisitor()) session_id = uuid.uuid4() - for visitor, visited_ch, handler in self._find_handlers(destination, visitors): - if self._handler_min_idle_time(handler): - reclaimed = await self._check_pel( - cmd=cmd, - handler=handler, - session_id=session_id, - ) - if reclaimed is not None: - return reclaimed - continue - + for visitor, visited_ch, handler in self._find_handlers( + destination=destination, + visitors=visitors, + cmd=cmd, + session_id=session_id, + ): msg = visitor.get_message( visited_ch, body, handler, ) - await self._put_pel(msg=msg, cmd=cmd, session_id=session_id, handler=handler) + self._put_pel(msg=msg, cmd=cmd, handler=handler, session_id=session_id) with anyio.fail_after(cmd.timeout): return await self._execute_handler(msg, handler, session_id=session_id) @@ -342,19 +339,18 @@ async def publish_batch(self, cmd: "RedisPublishCommand") -> int: for visitor, visited_ch, handler in self._find_handlers( {"list": cmd.destination}, (ListVisitor(),), + cmd=cmd, + session_id=session_id, ): casted_handler = cast("_ListHandlerMixin", handler) if casted_handler.list_sub.batch: - await self._check_pel(cmd=cmd, handler=handler, session_id=session_id) msg = visitor.get_message( visited_ch, data_to_send, casted_handler, ) - await self._put_pel( - msg=msg, cmd=cmd, session_id=session_id, handler=handler - ) + self._put_pel(msg=msg, cmd=cmd, handler=handler, session_id=session_id) await self._execute_handler(msg, handler, session_id=session_id) return 0 @@ -366,7 +362,7 @@ async def _execute_handler( session_id: uuid.UUID, ) -> "PubSubMessage": result = await handler.process_message(msg) - self._remove_pel(handler=handler, session_id=session_id, result=result) + self._remove_pel(handler=handler, result=result, session_id=session_id) return PubSubMessage( type="message", data=await build_message( @@ -385,30 +381,49 @@ def _find_handlers( self, destination: "_DestinationKwargs", visitors: "Sequence[Visitor]", + cmd: "RedisPublishCommand", + session_id: uuid.UUID, ) -> "Iterator[tuple[Visitor, str, LogicSubscriber]]": published_groups: set[tuple[str, str]] = set() - claimers: list[tuple[Visitor, str, LogicSubscriber]] = [] for handler in self.subscribers: # pragma: no branch for visitor in visitors: visited_ch = visitor.visit(**destination, sub=handler) if visited_ch is None: continue - - if isinstance(handler, _StreamHandlerMixin) and handler.stream_sub.group: - if handler.stream_sub.min_idle_time: - claimers.append((visitor, visited_ch, handler)) - break - - group_key = (visited_ch, handler.stream_sub.group) - if group_key in published_groups: - break - published_groups.add(group_key) - + if not self._return_handlers( + handler=handler, + visited_ch=visited_ch, + published_groups=published_groups, + cmd=cmd, + session_id=session_id, + ): + break yield visitor, visited_ch, handler break - yield from claimers + def _return_handlers( + self, + handler: "LogicSubscriber", + visited_ch: str, + published_groups: set[tuple[str, str]], + cmd: "RedisPublishCommand", + session_id: uuid.UUID, + ) -> bool: + if isinstance(handler, _StreamHandlerMixin) and handler.stream_sub.group: + group_key = (visited_ch, handler.stream_sub.group) + + if self._handler_min_idle_time(handler) and self._check_pel( + handler=handler, + cmd=cmd, + session_id=session_id, + ): + return True + if group_key in published_groups: + return False + published_groups.add(group_key) + return True + return True def _handler_group(self, handler: "LogicSubscriber") -> str | None: if isinstance(handler, _StreamHandlerMixin): @@ -423,25 +438,21 @@ def _handler_min_idle_time(self, handler: "LogicSubscriber") -> int | None: return handler.stream_sub.min_idle_time return None - async def _check_pel( + def _check_pel( self, handler: "LogicSubscriber", cmd: "RedisPublishCommand", session_id: uuid.UUID, - ) -> Optional["PubSubMessage"]: - pel_entry = self.pel.get_entry( + ) -> Optional["Entry"]: + return self.pel.get_entry( correlation_id=( cmd.correlation_id, session_id, self._handler_group(handler), ) ) - if pel_entry is None: - return None - - return await self._execute_handler(pel_entry.msg, handler, session_id=session_id) - async def _put_pel( + def _put_pel( self, handler: "LogicSubscriber", msg: Any, diff --git a/tests/brokers/redis/test_stream_group_pel.py b/tests/brokers/redis/test_stream_group_pel.py index 43cf5f769a..53fbd2e20a 100644 --- a/tests/brokers/redis/test_stream_group_pel.py +++ b/tests/brokers/redis/test_stream_group_pel.py @@ -22,11 +22,7 @@ async def worker_b(msg: str) -> None: ... pel = PEL() async with TestRedisBroker(broker, pel=pel) as br: await br.publish("data", stream="tasks") - - worker_a.mock.assert_called_once_with("data") - # group-b is unaffected by group-a's nack: it still gets delivered normally. - worker_b.mock.assert_called_once_with("data") - + assert len(pel._entries) == 1 @pytest.mark.redis() @@ -43,40 +39,61 @@ async def worker(msg: str) -> None: ... await br.publish("data", stream="tasks") worker.mock.assert_called_once_with("data") - # no_ack consumers are never tracked in the PEL at all. put_mock.assert_not_called() assert pel._entries == {} @pytest.mark.redis() @pytest.mark.asyncio() -async def test_only_min_idle_time_subscriber_processes_pel() -> None: +async def test_no_ack_policy_skips_pel_tracking_on_failure() -> None: broker = RedisBroker() - @broker.subscriber(stream=StreamSub("tasks", group="workers", consumer="w1")) + @broker.subscriber(stream=StreamSub("tasks", no_ack=True)) async def worker(msg: str) -> None: - raise NackMessage - - @broker.subscriber( - stream=StreamSub( - "tasks", - group="workers", - consumer="claimer", - min_idle_time=10000, - ), - ) - async def claimer(msg: str) -> None: ... + error_msg = "boom" + raise ValueError(error_msg) pel = PEL() async with TestRedisBroker(broker, pel=pel) as br: - with patch.object(pel, "get_entry", wraps=pel.get_entry) as get_entry_mock: + with ( + patch.object(pel, "put") as put_mock, + pytest.raises(ValueError, match="boom"), + ): await br.publish("data", stream="tasks") worker.mock.assert_called_once_with("data") - # worker nacked, so the entry stayed in the PEL and was reclaimed - claimer.mock.assert_called_once_with("data") - # only the min_idle_time consumer ever looks at the PEL, not `worker` - get_entry_mock.assert_called_once() + put_mock.assert_not_called() + assert pel._entries == {} + + +@pytest.mark.redis() +@pytest.mark.asyncio() +async def test_only_min_idle_time_subscriber_processes_pel() -> None: + call_order: list[str] = [] + while "worker" not in call_order: + broker = RedisBroker() + call_order = [] + + @broker.subscriber(stream=StreamSub("tasks", group="workers", consumer="w1")) + async def worker(msg: str) -> None: + call_order.append("worker") # noqa: B023 + raise NackMessage + + @broker.subscriber( + stream=StreamSub( + "tasks", + group="workers", + consumer="claimer", + min_idle_time=10000, + ), + ) + async def claimer(msg: str) -> None: + call_order.append("claimer") # noqa: B023 + + pel = PEL() + async with TestRedisBroker(broker, pel=pel) as br: + await br.publish("data", stream="tasks") + assert len(pel._entries) == 0 @pytest.mark.redis() @@ -120,5 +137,4 @@ async def worker(msg: str) -> None: await br.publish("data", stream="tasks") worker.mock.assert_called_once_with("data") - # nothing exists to reclaim it, so the pending entry just stays put. assert len(pel._entries) == 1 diff --git a/tests/docs/redis/stream/test_pel.py b/tests/docs/redis/stream/test_pel.py new file mode 100644 index 0000000000..c81c76fa6e --- /dev/null +++ b/tests/docs/redis/stream/test_pel.py @@ -0,0 +1,31 @@ +import pytest + + +@pytest.mark.redis() +@pytest.mark.asyncio() +async def test_pending_message_stays_without_a_claimer() -> None: + from docs.docs_src.redis.stream.pel_testing import ( + test_pending_message_stays_without_a_claimer as run, + ) + + await run() + + +@pytest.mark.redis() +@pytest.mark.asyncio() +async def test_each_group_gets_its_own_pel_entry() -> None: + from docs.docs_src.redis.stream.pel_testing import ( + test_each_group_gets_its_own_pel_entry as run, + ) + + await run() + + +@pytest.mark.redis() +@pytest.mark.asyncio() +async def test_no_ack_worker_never_touches_the_pel() -> None: + from docs.docs_src.redis.stream.pel_testing import ( + test_no_ack_worker_never_touches_the_pel as run, + ) + + await run() From 2f90561c87eecac7fb6679878d8d3c1d7f96c2e1 Mon Sep 17 00:00:00 2001 From: Apus Berliozi Date: Thu, 30 Jul 2026 23:04:35 +0700 Subject: [PATCH 10/10] feat(test-client redis):[#2927] fixed docs --- docs/docs/en/redis/streams/testing.md | 2 -- 1 file changed, 2 deletions(-) diff --git a/docs/docs/en/redis/streams/testing.md b/docs/docs/en/redis/streams/testing.md index 191794ee2b..4962581a31 100644 --- a/docs/docs/en/redis/streams/testing.md +++ b/docs/docs/en/redis/streams/testing.md @@ -94,5 +94,3 @@ async with TestRedisBroker(broker, pel=pel) as br: put_mock.assert_not_called() ``` - -Sharing a single `PEL` instance is also how you'd simulate two independent app instances competing for the same consumer group in a test - pass the same `pel` to multiple `TestRedisBroker` context managers instead of letting each create its own.