error_observer

Observer reporting the failures a pipeline runs into.

Errors travel upstream from the processor that raised them, and not all of them reach the end of that journey: a processor that answers for a failure itself — a service switcher that fails over to its next service, say — stops the error there. This observer reads each error where it is raised, so a session’s failure history holds the ones that were recovered from as well as the ones that surfaced.

class pipecat.observers.error_observer.ErrorEvent(*, message: str, category: ErrorCategory, exception_type: str | None = None, processor: str, processor_usable: bool, timestamp: float)[source]

Bases: BaseModel

One failure, as the processor that raised it described it.

Parameters:
  • message – What went wrong, in the words of the processor that failed.

  • category – Why it failed, drawn from ErrorCategory and independent of the provider that failed: rejected credentials, an unreachable service, a malformed request and so on.

  • exception_type – The name of the exception behind the failure, where one caused it. Failures group by this where a message, carrying the particulars of a single occurrence, is too specific to group by.

  • processor – The name of the processor that raised the error.

  • processor_usable – Whether that processor can still do its job. A processor that can’t keeps failing for as long as it is given work, so this separates a bad minute from the end of a capability.

  • timestamp – Unix timestamp of the failure.

class pipecat.observers.error_observer.ErrorObserver(*, time_source: Callable[[], float]=<built-in function time>, **kwargs)[source]

Bases: BaseObserver

Reports each error a pipeline raises, once, where it is raised.

An error is reported at its origin rather than where it ends up, and named for the processor that raised it rather than the one that passed it along.

Events:
on_error(observer, event): Emitted for each error, as an

ErrorEvent.

Example:

observer = ErrorObserver()

@observer.event_handler("on_error")
async def on_error(observer, event):
    logger.info(event.model_dump_json())
__init__(*, time_source: Callable[[], float]=<built-in function time>, **kwargs)[source]

Initialize the error observer.

Parameters:
  • time_source – Reads the current time in seconds. Supplying one lets a test place failures without waiting.

  • **kwargs – Additional arguments passed to parent class.

async on_push_frame(data: FramePushed)[source]

Report an error frame.

The first push of an error comes from the processor that failed.

Parameters:

data – Frame push event containing the frame and direction.