Skip to content

feat: add subscriber and broker exception handlers - #3243

Draft
DMINQ wants to merge 7 commits into
ag2ai:mainfrom
DMINQ:feat/2945-exception-handlers
Draft

DMINQ wants to merge 7 commits into
ag2ai:mainfrom
DMINQ:feat/2945-exception-handlers

Conversation

@DMINQ

@DMINQ DMINQ commented Sep 26, 2026

Copy link
Copy Markdown

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

  • Four automated test cases pass, covering handler precedence,
    fallback results and dispatch through consume().
  • The consume() test was verified to fail when the dispatch call
    was removed, and to pass with it restored.
  • Manually verified with a real Kafka broker: local handling,
    broker fallback, unhandled exceptions and continued processing.
  • Ruff lint and formatting checks passed for the changed files.
  • Mypy, Pyright and Pyrefly passed for the initial implementation.
  • Targeted Mypy checks passed after adding the consume() test.

Remaining work

  • Cover missing handlers and remaining sync/async combinations.
  • Integrate background-task error handling.
  • Expose the API across the remaining brokers.
  • Add documentation and runnable examples.
  • Complete the remaining project checks.

Type of change

  • New feature
  • This change requires a documentation update

Checklist

  • Self-reviewed the current changes.
  • Added automated tests for the implemented behavior.
  • Added documentation and tested examples.
  • Passed the full project lint checks (just lint).
  • Passed the relevant coverage checks.
  • Passed the full static analysis suite (just static-analysis).

@CLAassistant

CLAassistant commented Sep 26, 2026 •

Copy link
Copy Markdown

CLA assistant check
All committers have signed the CLA.

@github-actions github-actions Bot added the AioKafka Issues related to `faststream.kafka` module label Sep 26, 2026
@github-actions github-actions Bot added the Confluent Issues related to `faststream.confluent` module label Sep 28, 2026

@IvanKirpichnikov IvanKirpichnikov left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 IvanKirpichnikov left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread faststream/_internal/configs/broker.py Outdated
Comment on lines +37 to +42
broker_exception_handler: Optional["ExceptionHandler"] = None
_broker_exception_handler: Optional["AsyncExceptionHandler"] = field(
default=None,
init=False,
repr=False,
)

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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(...).

Comment thread faststream/_internal/broker/broker.py Outdated
Comment on lines +53 to +61
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,
)

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Let’s do everything the same as we did with BrokerConfig

Comment on lines +92 to +101
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,
)

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

then it is just self._exception_handler = config.exception_hadnler

Comment thread faststream/kafka/broker/broker.py Outdated
broker_decoder=decoder,
broker_codec=codec,
broker_parser=parser,
broker_exception_handler=exception_handler,

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Here we’ll wrap it in to_async I mentioned this above

Comment on lines +36 to +38
if msg == "hello":
raise error
event2.set()

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I don’t really understand the test is complicated. Just do raise error, why two conditions?

Comment thread tests/brokers/base/exception_handlers.py
Comment thread tests/brokers/kafka/test_exception_handlers.py
Comment thread tests/brokers/base/exception_handlers.py

@IvanKirpichnikov IvanKirpichnikov left a comment •

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

And add docs

Comment thread faststream/app.py
Comment on lines +78 to +84
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,
)

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggested change
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

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment on lines +345 to +361
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)

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment on lines +342 to +363
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

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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

Comment on lines +386 to +392
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

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

and then this code will turn into

result = await self._handle_exception(exc)
if not result:
    raise

Comment on lines +77 to +79
exception_handler=(
to_async(exception_handler) if exception_handler is not None else None
),

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

And add some kind of ensure_exception_handler in faststream/_internal/utils/functions.py that will do this internally.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

AioKafka Issues related to `faststream.kafka` module Confluent Issues related to `faststream.confluent` module

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants