Repository navigation
Conversation
IvanKirpichnikov
left a comment
There was a problem hiding this comment.
Hi. Sorry it took so long.
The idea of using to_async in __post_init__ didn’t really work for me. It ended up feeling a bit cumbersome. Let’s go back to your approach.
Second, why are you wrapping it in _build_fastdepends_model? Let’s do it in __init__ instead. In theory, nothing should break.
Third, the tests.
Let’s write e2e tests. The current tests are fragile because you monkeypatch a lot of things and call private methods.
I think it would be better to write e2e tests like this:
async def test_subscriber_exception_handler_true_suppresses_error(
self,
queue: str,
event: Event,
mock: MagicMock,
) -> None:
error = RuntimeError()
async def exception_handler(exc: BaseException) -> bool:
mock(exc)
return True
broker = self.get_broker()
args, kwargs = self.get_subscriber_params(queue, exception_handler=exception_handler)
@broker.subscriber(*args, **kwargs)
async def handler(msg: Any) -> NoReturn:
event.set()
raise error
async with self.patch_broker(broker) as br:
await asyncio.wait(
(
asyncio.create_task(br.publish(queue, None)),
asyncio.create_task(event.wait()),
),
timeout=self.timeout,
)
mock.assert_called_one_with(error)
IvanKirpichnikov
left a comment
There was a problem hiding this comment.
We also need to add another global exception handler at the FastStream level, which will also be taken into account when handling errors, as the very last one. Tests will need to be added for it.
| broker_exception_handler: Optional["ExceptionHandler"] = None | ||
| _broker_exception_handler: Optional["AsyncExceptionHandler"] = field( | ||
| default=None, | ||
| init=False, | ||
| repr=False, | ||
| ) |
There was a problem hiding this comment.
So, in BrokerConfig, keep a single field, broker_exception_handler: Optional["AsyncExceptionHandler"], but allow broker_exception_handler: Optional["ExceptionHandler"] to be passed directly to the brokers. In other words, when creating BrokerConfig, we would pass BrokerConfig(..., broker_exception_handler=to_async(broker_exception_handler)) and then in __post_init__, let's decorate it using apply_types self.broker_exception_handler = apply_types(...).
| if config.broker_exception_handler is not None: | ||
| async_handler: AsyncExceptionHandler = to_async( | ||
| config.broker_exception_handler, | ||
| ) | ||
| config._broker_exception_handler = apply_types( | ||
| async_handler, | ||
| serializer_cls=config.fd_config._serializer, | ||
| context__=config.context, | ||
| ) |
There was a problem hiding this comment.
then it is just self._exception_handler = config.exception_hadnler
| @dataclass(kw_only=True) | ||
| class SubscriberUsecaseConfig(EndpointConfig): | ||
| no_reply: bool = False | ||
| exception_handler: "ExceptionHandler | None" = None |
There was a problem hiding this comment.
Let’s do everything the same as we did with BrokerConfig
| self._exception_handler_call: AsyncExceptionHandler | None = None | ||
| if config.exception_handler is not None: | ||
| async_handler: AsyncExceptionHandler = to_async( | ||
| config.exception_handler, | ||
| ) | ||
| self._exception_handler_call = apply_types( | ||
| async_handler, | ||
| serializer_cls=self._outer_config.fd_config._serializer, | ||
| context__=self._outer_config.context, | ||
| ) |
There was a problem hiding this comment.
then it is just self._exception_handler = config.exception_hadnler
| broker_decoder=decoder, | ||
| broker_codec=codec, | ||
| broker_parser=parser, | ||
| broker_exception_handler=exception_handler, |
There was a problem hiding this comment.
Here we’ll wrap it in to_async I mentioned this above
| if msg == "hello": | ||
| raise error | ||
| event2.set() |
There was a problem hiding this comment.
I don’t really understand the test is complicated. Just do raise error, why two conditions?
| self._exception_handler: AsyncExceptionHandler | None = None | ||
| if exception_handler is not None: | ||
| self._exception_handler = apply_types( | ||
| to_async(exception_handler), | ||
| serializer_cls=self.config._serializer, | ||
| context__=self.context, | ||
| ) |
There was a problem hiding this comment.
| self._exception_handler: AsyncExceptionHandler | None = None | |
| if exception_handler is not None: | |
| self._exception_handler = apply_types( | |
| to_async(exception_handler), | |
| serializer_cls=self.config._serializer, | |
| context__=self.context, | |
| ) | |
| if exception_handler is not None: | |
| self._exception_handler = apply_types( | |
| to_async(exception_handler), | |
| serializer_cls=self.config._serializer, | |
| ) | |
| else: | |
| self._exception_handler = None |
There was a problem hiding this comment.
Why did you create get_consume_message and call subscriber.consume?
Just use await broker.publish(None, queue). When your broker is in test mode, this code will wait for the handler to execute.
| local_handler = cast( | ||
| "Callable[..., Awaitable[bool]] | None", | ||
| self._exception_handler, | ||
| ) | ||
| if local_handler is not None and await local_handler(exc, context__=context): | ||
| return True | ||
|
|
||
| broker_handler = cast( | ||
| "Callable[..., Awaitable[bool]] | None", | ||
| self._outer_config.broker_exception_handler, | ||
| ) | ||
| if broker_handler is not None and await broker_handler(exc, context__=context): | ||
| return True | ||
|
|
||
| app_handler = self._app_exception_handler | ||
| if app_handler is not None: | ||
| return await app_handler(exc) |
There was a problem hiding this comment.
typing.cast will not be needed if you remove the context__ argument passing. Passing it is unnecessary, as you are already passing it to apply_types.
| async def _handle_exception(self, exc: BaseException) -> bool: | ||
| context = self._outer_config.context | ||
|
|
||
| local_handler = cast( | ||
| "Callable[..., Awaitable[bool]] | None", | ||
| self._exception_handler, | ||
| ) | ||
| if local_handler is not None and await local_handler(exc, context__=context): | ||
| return True | ||
|
|
||
| broker_handler = cast( | ||
| "Callable[..., Awaitable[bool]] | None", | ||
| self._outer_config.broker_exception_handler, | ||
| ) | ||
| if broker_handler is not None and await broker_handler(exc, context__=context): | ||
| return True | ||
|
|
||
| app_handler = self._app_exception_handler | ||
| if app_handler is not None: | ||
| return await app_handler(exc) | ||
|
|
||
| return False |
There was a problem hiding this comment.
By the way, this entire function can be represented as a “Chain of Responsibility” pattern.
Do something like
handlers = [...]
for handler in handlers:
if handler:
res = handler()
if res:
return True
return False| has_handler = ( | ||
| self._exception_handler is not None | ||
| or self._outer_config.broker_exception_handler is not None | ||
| or self._app_exception_handler is not None | ||
| ) | ||
| if has_handler and not await self._handle_exception(exc): | ||
| raise |
There was a problem hiding this comment.
and then this code will turn into
result = await self._handle_exception(exc)
if not result:
raise| exception_handler=( | ||
| to_async(exception_handler) if exception_handler is not None else None | ||
| ), |
There was a problem hiding this comment.
And add some kind of ensure_exception_handler in faststream/_internal/utils/functions.py that will do this internally.
Description
Adds configurable exception handlers following the proposal in:
#3095 (comment)
The subscriber handler runs first. If it is absent or returns False,
the broker handler runs. Both synchronous and asynchronous handlers
are supported.
This draft currently exposes the API for Kafka and adds shared
dispatch logic to SubscriberUsecase.consume(). Unhandled exceptions
retain the existing consume() behavior. Middleware logging and
acknowledgement behavior remain unchanged.
Related to #2945. Background-task error handling is not implemented
yet, so this draft does not resolve the original Redis scenario.
Validation
fallback results and dispatch through consume().
was removed, and to pass with it restored.
broker fallback, unhandled exceptions and continued processing.
Remaining work
Type of change
Checklist
just lint).just static-analysis).