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?
There was a problem hiding this comment.
Let’s do the tests in memory. There’s no need for a real connection to the brokers.
There was a problem hiding this comment.
Why are these tests separated separately for kafka? And why is subscriber.tasks used in them? There’s no need for that.
There was a problem hiding this comment.
It would also be worth adding tests for multi‑brokers. We need to check that the exception handler for a specific broker is called, and not someone else’s.
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).