Skip to content

feat: add subscriber and broker exception handlers - #3243

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

DMINQ wants to merge 5 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 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 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

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?

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 the tests in memory. There’s no need for a real connection to the brokers.

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 are these tests separated separately for kafka? And why is subscriber.tasks used in them? There’s no need for that.

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.

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.

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