Skip to content

feat: Allow capturing exception on task execution - #3095

Open
pbonneaudiabolocom wants to merge 3 commits into
ag2ai:mainfrom
pbonneaudiabolocom:handle-consume-error-in-subscriber
Open

pbonneaudiabolocom wants to merge 3 commits into
ag2ai:mainfrom
pbonneaudiabolocom:handle-consume-error-in-subscriber

Conversation

@pbonneaudiabolocom

@pbonneaudiabolocom pbonneaudiabolocom commented Sep 4, 2026 •

Copy link
Copy Markdown

Description

It allows to customize the task exception managment + the consume error behaviour

Fixes #2945

Type of change

Please delete options that are not relevant.

  • Documentation (typos, code examples, or any documentation updates)
  • Bug fix (a non-breaking change that resolves an issue)
  • New feature (a non-breaking change that adds functionality)
  • Breaking change (a fix or feature that would disrupt existing functionality)
  • This change requires a documentation update

Checklist

  • My code adheres to the style guidelines of this project (just lint shows no errors)
  • I have conducted a self-review of my own code
  • I have made the necessary changes to the documentation
  • My changes do not generate any new warnings
  • I have added tests to validate the effectiveness of my fix or the functionality of my new feature
  • Both new and existing unit tests pass successfully on my local environment by running just test-coverage
  • I have ensured that static analysis tests are passing by running just static-analysis
  • I have included code examples to illustrate the modifications

…ne to overload the task exception that could be generated.
@github-actions github-actions Bot added documentation Improvements or additions to documentation Redis Issues related to `faststream.redis` module and Redis features labels Sep 4, 2026
@pbonneaudiabolocom

Copy link
Copy Markdown
Author

Pursuing efforts of PR #2947

@IvanKirpichnikov
@powersemmi
@Lancetnik

@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. Let’s do it this way: let the subscriber accept exception_handler: Callable[[BaseException], bool].

If exception_handler returns True, the exception is suppressed. If it returns False, the exception is considered unhandled, following the same logic as context managers.

This means that this hook should be available on every subscriber. Not only for redis

@IvanKirpichnikov

Copy link
Copy Markdown
Collaborator

I don’t like this approach for several reasons:

Before defining your own handler, you have to save the “old” one:

  • default_task_exception_handler = subscriber.handle_task_exception
  • You have to override a method on the subscriber. Not a fan of that.
  • It accepts all sorts of func, func_args, and func_kwargs and uses add_task.

@pbonneaudiabolocom

Copy link
Copy Markdown
Author

I agree, the issue is that you can't overload subscriber because it's provided by the broker.
And you do not instantiate it either.

So, passing another argument, exception_handler: Callable[[BaseException], bool], is going to be as difficult as switching a method to replace it by another.

After, my knowledge on how faststream works is limited. I'm a user, not so knowledgable about what's under the hood.

The only object we do have access to is the broker, do we really want the broker to be the entry point for a subscriber exception manager ?

I'm honnestly trying to find a way, but I find many average solution, but no good one when I'm satisfied by the result.

@IvanKirpichnikov

Copy link
Copy Markdown
Collaborator

The only object we do have access to is the broker, do we really want the broker to be the entry point for a subscriber exception manager ?

I mean passing the exception_handler to the subscriber.

Although, if we take the idea further, I think we could also support setting the exception_handler at the broker level, but the question arises as to what order to call them in.
I think we should first call the exception_handler on the subscriber, and if it returns False, then call the exception_handler on the broker.

@pbonneaudiabolocom

Copy link
Copy Markdown
Author

I still don't get how we are supposed to access the subscriber to do that.
This object is created by the broker, so you really want to do something such as

broker.subscriber.add_exception_handler()

I mean, it could be done for sure, but then you'll come back quickly to something such as getting the default exception_handler for the subscriber, and replacing it by a function taht do some stuff and call the default one in case needed.

I cannot figure what you want to do...

@IvanKirpichnikov

Copy link
Copy Markdown
Collaborator

I still don't get how we are supposed to access the subscriber to do that. This object is created by the broker, so you really want to do something such as

broker.subscriber.add_exception_handler()

I mean, it could be done for sure, but then you'll come back quickly to something such as getting the default exception_handler for the subscriber, and replacing it by a function taht do some stuff and call the default one in case needed.

I cannot figure what you want to do...

I mean something like that.

def global_exception_handler(exc: BaseException) -> bool:
    if isinstance(exc, GlobalError):
        ...
        return True
    return False

broker = NatsBroker(..., exception_handler=global_exception_handler)


def my_exception_handler(exc: BaseException) -> bool:
    if isinstance(exc, MyError):
        # do custom logic
        ...

        # i'm process error
        return True

    # i'm not process error
    return False

@broker.subscriber("subject", exception_handler=my_exception_handler)
async def handle(body: Body):
    if body.act == 1:
        raise MyError
    else:
        raise GlobalError

@pbonneaudiabolocom

Copy link
Copy Markdown
Author

ok, so on top of the subscriber exception_handler, you want also a broker exception handler.
Both being injected using constructor.

That could be elegant.

Are we allowing to overide the default exception handling, or to call it ?
With such a system, how could we implement the logic :

If groupError:
Do something custom
Else:
Use regular excetion handling //could be also disable exeption by default.

Because I think that's going to be a common usecase. Target an error, do a specific action, but let faststream defaut behaviour in any other usecases.
Implementing the custom_exception_handler outside the subscriber(and before it's creation), we will not have access to the subscriber, nor it's values.

@IvanKirpichnikov

IvanKirpichnikov commented Sep 6, 2026 •

Copy link
Copy Markdown
Collaborator

ok, so on top of the subscriber exception_handler, you want also a broker exception handler. Both being injected using constructor.

That could be elegant.

Are we allowing to overide the default exception handling, or to call it ? With such a system, how could we implement the logic :

If groupError:
Do something custom
Else:
Use regular excetion handling //could be also disable exeption by default.

Because I think that's going to be a common usecase. Target an error, do a specific action, but let faststream defaut behaviour in any other usecases. Implementing the custom_exception_handler outside the subscriber(and before it's creation), we will not have access to the subscriber, nor it's values.

My idea is as follows.

First, the exception_handler is called for the subscriber. If it returns False, then the exception_handler for the broker is executed. If it also returns False, the default error handling is triggered.

Is everything clear to you? If not, I’m ready to answer your questions.

@pbonneaudiabolocom

Copy link
Copy Markdown
Author

Hi,

Your idea seems nice, but I'm not sure I'm the right one to implement it. I barely understand faststream logic, so except burning token, my value will be very low here...

I don't know how I could assert anything... especially with such a large scope impact.

@IvanKirpichnikov

Copy link
Copy Markdown
Collaborator

It doesn’t seem too difficult. It’s worth a try.

Patience and hard work will conquer all :)

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

documentation Improvements or additions to documentation Redis Issues related to `faststream.redis` module and Redis features

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Bug: group removal handling in Redis Stream not functional

2 participants