Source code for pipecat.observers.base_observer

#
# Copyright (c) 2024-2026, Daily
#
# SPDX-License-Identifier: BSD 2-Clause License
#

"""Base observer classes for monitoring frame flow in the Pipecat pipeline.

This module provides the foundation for observing frame transfers between
processors without modifying the pipeline structure. Observers can be used
for logging, debugging, analytics, and monitoring pipeline behavior.
"""

from dataclasses import dataclass
from typing import TYPE_CHECKING

from pipecat.frames.frames import Frame
from pipecat.utils.base_object import BaseObject
from pipecat.utils.deprecation import deprecated

if TYPE_CHECKING:
    from pipecat.processors.frame_processor import FrameDirection, FrameProcessor


[docs] @dataclass class FrameProcessed: """Event data for frame processing in the pipeline. Represents an event where a frame is being processed by a processor. This data structure is typically used by observers to track the flow of frames through the pipeline for logging, debugging, or analytics purposes. Parameters: processor: The processor processing the frame. frame: The frame being processed. direction: The direction of the frame (e.g., downstream or upstream). timestamp: The time when the frame was pushed, based on the pipeline clock. """ processor: "FrameProcessor" frame: Frame direction: "FrameDirection" timestamp: int
[docs] @dataclass class FramePushed: """Event data for frame transfers between processors in the pipeline. Represents an event where a frame is pushed from one processor to another within the pipeline. This data structure is typically used by observers to track the flow of frames through the pipeline for logging, debugging, or analytics purposes. Parameters: source: The processor sending the frame. destination: The processor receiving the frame. frame: The frame being transferred. direction: The direction of the transfer (e.g., downstream or upstream). timestamp: The time when the frame was pushed, based on the pipeline clock. first_push: Whether this is the first time the frame is pushed. A frame is pushed again by every processor that passes it along. """ source: "FrameProcessor" destination: "FrameProcessor" frame: Frame direction: "FrameDirection" timestamp: int first_push: bool = True
[docs] @dataclass class ProcessorSetUp: """Event data for a processor having been set up. Processors are set up concurrently and before any frame flows, so this is what a timing observer measures the work a processor does to get ready by. The times come from :func:`time.monotonic_ns`, since the pipeline clock only starts once the pipeline does. Parameters: processor: The processor that was set up. started_at_ns: When the processor's ``setup()`` began. finished_at_ns: When the processor's ``setup()`` returned. """ processor: "FrameProcessor" started_at_ns: int finished_at_ns: int
[docs] @deprecated( "`StartupWarmup` is deprecated since 1.12.0 and will be removed in 2.0.0. No replacement." ) @dataclass class StartupWarmup: """Event data for the framework having warmed its deferred imports. .. deprecated:: 1.12.0 No replacement. Nothing warms deferred imports at startup, so this event is never emitted. Will be removed in 2.0.0. Parameters: started_at_ns: When warming began. finished_at_ns: When warming finished. """ started_at_ns: int finished_at_ns: int
[docs] class BaseObserver(BaseObject): """Base class for pipeline frame observers. Observers can view all frames that flow through the pipeline without needing to inject processors into the pipeline structure. This enables non-intrusive monitoring capabilities such as frame logging, debugging, performance analysis, and analytics collection. A frame is pushed again by every processor that passes it along, and by default an observer observes every push. An observer that handles a frame once, such as one that reports the moment a frame represents, is created with ``observe_every_push=False`` and is told about a frame only when it is first pushed. """
[docs] def __init__(self, *, observe_every_push: bool = True, **kwargs): """Initialize the observer. Args: observe_every_push: Whether to observe every push of a frame, rather than only its first. Defaults to True. **kwargs: Additional arguments passed to the parent class. """ super().__init__(**kwargs) self._observe_every_push = observe_every_push
@property def observe_every_push(self) -> bool: """Whether the observer observes every push of a frame, not only the first.""" return self._observe_every_push
[docs] async def on_process_frame(self, data: FrameProcessed): """Handle the event when a frame is being processed by a processor. This method should be implemented by subclasses to define specific behavior (e.g., logging, monitoring, debugging) when a frame is being processed by a processor. Args: data: The event data containing details about the frame processing. """ pass
[docs] async def on_push_frame(self, data: FramePushed): """Handle the event when a frame is pushed from one processor to another. This method should be implemented by subclasses to define specific behavior (e.g., logging, monitoring, debugging) when a frame is transferred through the pipeline. A frame is pushed again by every processor that passes it along, so this is called once per hop, with ``data.first_push`` set on the first one. An observer created with ``observe_every_push=False`` is only called for that first push. Args: data: The event data containing details about the frame transfer. """ pass
[docs] async def on_processor_setup(self, data: ProcessorSetUp): """Handle the event when a processor has been set up. A processor connects and does its other slow start-up work here, so this is where that cost can be measured. Processors are set up concurrently, so these arrive in the order they finish rather than in pipeline order. Args: data: The event data containing details about the processor setup. """ pass
[docs] @deprecated( "`BaseObserver.on_startup_warmup` is deprecated since 1.12.0 and will be removed in " "2.0.0. No replacement." ) async def on_startup_warmup(self, data: StartupWarmup): """Handle the event when the framework has warmed its deferred imports. .. deprecated:: 1.12.0 No replacement. Nothing warms deferred imports at startup, so this is never called. Will be removed in 2.0.0. Args: data: The event data containing details about the warming. """ pass
[docs] async def on_pipeline_started(self): """Called when the pipeline has fully started. Fired after the ``StartFrame`` has been processed by all processors in the pipeline, including nested ``ParallelPipeline`` branches. """ pass