worker

Pipeline worker implementation for managing frame processing pipelines.

This module provides the main PipelineWorker class that orchestrates pipeline execution, frame routing, lifecycle management, and monitoring capabilities including heartbeats, idle detection, and observer integration.

class pipecat.pipeline.worker.IdleFrameObserver(*, idle_event: Event, idle_timeout_frames: tuple[type[Frame], ...])[source]

Bases: BaseObserver

Idle timeout observer.

This observer waits for specific frames being generated in the pipeline. If the frames are generated the given asyncio event is set. If the event is not set it means the pipeline is probably idle.

__init__(*, idle_event: Event, idle_timeout_frames: tuple[type[Frame], ...])[source]

Initialize the observer.

Parameters:
  • idle_event – The event to set if the idle timeout frames are being pushed.

  • idle_timeout_frames – A tuple with the frames that should set the event when received

async on_push_frame(data: FramePushed)[source]

Callback executed when a frame is pushed in the pipeline.

Parameters:

data – The frame push event data.

class pipecat.pipeline.worker.ProcessorUnusablePolicy(*values)[source]

Bases: Enum

What a pipeline worker does when a processor can no longer do its job.

An unusable processor keeps failing for as long as the pipeline keeps using it, so the pipeline has to decide whether it is still worth running. A processor becomes unusable through a permanent error category or through FrameProcessor.push_error() with force_treat_as_permanent=True.

Parameters:
  • CONTINUE – Report the error and keep running. The application can decide what to do, for example by failing over to another provider.

  • END – End the pipeline gracefully, letting queued frames drain first.

  • CANCEL – Cancel the pipeline immediately, abandoning queued frames.

CONTINUE = 'continue'
END = 'end'
CANCEL = 'cancel'
class pipecat.pipeline.worker.PipelineParams(*, audio_in_sample_rate: int = 16000, audio_out_sample_rate: int = 24000, enable_heartbeats: bool = False, enable_metrics: bool = False, enable_usage_metrics: bool = False, heartbeats_period_secs: float = 1.0, heartbeats_monitor_secs: float = 10.0, report_only_initial_ttfb: bool = False, send_initial_empty_metrics: bool = True, start_metadata: dict[str, ~typing.Any]=<factory>)[source]

Bases: BaseModel

Configuration parameters for pipeline execution.

These parameters are usually passed to all frame processors through StartFrame. For other generic pipeline worker parameters use PipelineWorker constructor arguments instead.

Parameters:
  • audio_in_sample_rate – Input audio sample rate in Hz.

  • audio_out_sample_rate – Output audio sample rate in Hz.

  • enable_heartbeats – Whether to enable heartbeat monitoring.

  • enable_metrics – Whether to enable metrics collection.

  • enable_usage_metrics – Whether to enable usage metrics.

  • heartbeats_period_secs – Period between heartbeats in seconds.

  • heartbeats_monitor_secs – Timeout (in seconds) before warning about missed heartbeats. Defaults to 10 seconds.

  • report_only_initial_ttfb – Whether to report only initial time to first byte.

  • send_initial_empty_metrics – Whether to send initial empty metrics.

  • start_metadata – Additional metadata for pipeline start.

class pipecat.pipeline.worker.PipelineWorker(pipeline: ~pipecat.pipeline.base_pipeline.BasePipeline, *, active: bool = True, additional_span_attributes: dict | None = None, app_resources: ~typing.Any = None, bridged: tuple[str, ...] | None = None, cancel_on_idle_timeout: bool = True, cancel_runner_on_idle_timeout: bool = True, cancel_timeout_secs: float = 20.0, check_dangling_tasks: bool = True, clock: ~pipecat.clocks.base_clock.BaseClock | None = None, conversation_id: str | None = None, enable_tracing: bool = False, enable_turn_tracking: bool = True, handle_flush_frame: bool | None = None, enable_rtvi: bool | None = None, exclude_frames: tuple[type[~pipecat.frames.frames.Frame], ...] | None = None, idle_timeout_frames: tuple[type[~pipecat.frames.frames.Frame], ...] = (<class 'pipecat.frames.frames.BotSpeakingFrame'>, <class 'pipecat.frames.frames.InterimTranscriptionFrame'>, <class 'pipecat.frames.frames.TranscriptionFrame'>, <class 'pipecat.frames.frames.UserSpeakingFrame'>, <class 'pipecat.frames.frames.UserStartedSpeakingFrame'>), idle_timeout_secs: float | None = 300, name: str | None = None, observers: list[~pipecat.observers.base_observer.BaseObserver] | None = None, processor_unusable_policy: ~pipecat.pipeline.worker.ProcessorUnusablePolicy = ProcessorUnusablePolicy.CONTINUE, params: ~pipecat.pipeline.worker.PipelineParams | None = None, rtvi_processor: ~pipecat.processors.frameworks.rtvi.processor.RTVIProcessor | None = None, rtvi_observer_params: ~pipecat.processors.frameworks.rtvi.observer.RTVIObserverParams | None = None, setup_timeout_secs: float = 20.0, start_timeout_secs: float = 20.0, task_manager: ~pipecat.utils.asyncio.task_manager.BaseTaskManager | None = None, tool_resources: ~typing.Any = None)[source]

Bases: BaseWorker

Manages the execution of a pipeline, handling frame processing and worker lifecycle.

This class orchestrates pipeline execution with comprehensive monitoring, event handling, and lifecycle management. It provides event handlers for various pipeline states and frame types, idle detection, heartbeat monitoring, and observer integration.

Event handlers available:

  • on_frame_reached_upstream: Called when upstream frames reach the source

  • on_frame_reached_downstream: Called when downstream frames reach the sink

  • on_heartbeat_timeout: Called when a heartbeat frame is not received within the monitor timeout.

    Fires repeatedly every heartbeats_monitor_secs for as long as the stall persists.

  • on_idle_timeout: Called when pipeline is idle beyond timeout threshold

  • on_pipeline_started: Called when pipeline starts with StartFrame

  • on_pipeline_finished: Called after the pipeline has reached any terminal state.

    This includes:

    • StopFrame: pipeline was stopped (processors keep connections open)

    • EndFrame: pipeline ended normally

    • CancelFrame: pipeline was cancelled

    Use this event for cleanup, logging, or post-processing tasks. Users can inspect the frame if they need to handle specific cases.

  • on_pipeline_timeout: Called when a frame the worker was waiting on never

    reached the end of the pipeline, meaning a processor is blocked. Inspect the frame to tell the cases apart:

    • StartFrame: the pipeline never started and is being torn down

    • CancelFrame: the pipeline was cancelled but did not drain

    There is no EndFrame case: ending flushes whatever is queued, so the worker waits for it however long that takes.

  • on_pipeline_error: Called when an error occurs with ErrorFrame. Handler can read

    frame.processor.is_usable to distinguish between errors that end a processor’s usability from those that don’t.

Example:

@worker.event_handler("on_frame_reached_upstream")
async def on_frame_reached_upstream(worker, frame):
    ...

@worker.event_handler("on_heartbeat_timeout")
async def on_pipeline_heartbeat_timeout(worker):
    ...

@worker.event_handler("on_idle_timeout")
async def on_pipeline_idle_timeout(worker):
    ...

@worker.event_handler("on_pipeline_started")
async def on_pipeline_started(worker, frame):
    ...

@worker.event_handler("on_pipeline_finished")
async def on_pipeline_finished(worker, frame):
    ...

@worker.event_handler("on_pipeline_timeout")
async def on_pipeline_timeout(worker, frame):
    ...

@worker.event_handler("on_pipeline_error")
async def on_pipeline_error(worker, frame):
    ...

@worker.event_handler("on_setup_timeout")
async def on_setup_timeout(worker):
    ...
__init__(pipeline: ~pipecat.pipeline.base_pipeline.BasePipeline, *, active: bool = True, additional_span_attributes: dict | None = None, app_resources: ~typing.Any = None, bridged: tuple[str, ...] | None = None, cancel_on_idle_timeout: bool = True, cancel_runner_on_idle_timeout: bool = True, cancel_timeout_secs: float = 20.0, check_dangling_tasks: bool = True, clock: ~pipecat.clocks.base_clock.BaseClock | None = None, conversation_id: str | None = None, enable_tracing: bool = False, enable_turn_tracking: bool = True, handle_flush_frame: bool | None = None, enable_rtvi: bool | None = None, exclude_frames: tuple[type[~pipecat.frames.frames.Frame], ...] | None = None, idle_timeout_frames: tuple[type[~pipecat.frames.frames.Frame], ...] = (<class 'pipecat.frames.frames.BotSpeakingFrame'>, <class 'pipecat.frames.frames.InterimTranscriptionFrame'>, <class 'pipecat.frames.frames.TranscriptionFrame'>, <class 'pipecat.frames.frames.UserSpeakingFrame'>, <class 'pipecat.frames.frames.UserStartedSpeakingFrame'>), idle_timeout_secs: float | None = 300, name: str | None = None, observers: list[~pipecat.observers.base_observer.BaseObserver] | None = None, processor_unusable_policy: ~pipecat.pipeline.worker.ProcessorUnusablePolicy = ProcessorUnusablePolicy.CONTINUE, params: ~pipecat.pipeline.worker.PipelineParams | None = None, rtvi_processor: ~pipecat.processors.frameworks.rtvi.processor.RTVIProcessor | None = None, rtvi_observer_params: ~pipecat.processors.frameworks.rtvi.observer.RTVIObserverParams | None = None, setup_timeout_secs: float = 20.0, start_timeout_secs: float = 20.0, task_manager: ~pipecat.utils.asyncio.task_manager.BaseTaskManager | None = None, tool_resources: ~typing.Any = None)[source]

Initialize the PipelineWorker.

Parameters:
  • pipeline – The pipeline to execute.

  • active – Whether the worker starts active. Forwarded to BaseWorker.

  • additional_span_attributes – Optional dictionary of attributes to propagate as OpenTelemetry conversation span attributes.

  • app_resources – Optional application-defined bag of anything your application code may want to share across this session (DB handles, HTTP clients, etc.), passed by reference. Pipecat passes it through untouched and exposes it on the worker itself as worker.app_resources and passes it to tool handlers as FunctionCallParams.app_resources. The framework never copies or clears this object; the caller retains their handle and can read any mutations after the worker finishes.

  • bridged – Bridge configuration. None means the pipeline is not bridged. An empty tuple () wraps the pipeline with bus edge processors that accept frames from all bridges. A tuple of names like ("voice",) accepts only frames from those bridges. The bus comes from attach() (called by the runner).

  • handle_flush_frame – Whether this worker answers a flush probe, bouncing it at the sink and completing it at the source. Defaults to whether the pipeline is unbridged, so a worker wired into someone else’s topology leaves the probe to travel on and be completed by the pipeline that owns the transport. A bridged worker with no such peer never completes a flush and every flush_pipeline() call waits out its timeout.

  • cancel_on_idle_timeout – Whether reaching the idle timeout should cancel the pipeline worker. When False, the idle event still fires on_idle_timeout but the worker is left alone (and cancel_runner_on_idle_timeout is ignored too: opting out of local cancellation also opts out of the runner-wide cancel).

  • cancel_runner_on_idle_timeout – When cancel_on_idle_timeout is also True, whether reaching the idle timeout should also cancel the entire WorkerRunner. The worker is always cancelled first; when this is True the worker also emits a BusCancelMessage so the runner broadcasts cancellation to every other root worker. Defaults to True so a multi-worker bot’s helpers shut down with the main pipeline; set to False for a sidecar PipelineWorker that should self-cancel on idle without bringing down its peers.

  • cancel_timeout_secs – Timeout (in seconds) to wait for cancellation to happen cleanly.

  • check_dangling_tasks – Whether to warn about tasks left running when the worker finishes. Only applies when the worker owns its task manager; otherwise the runner reports dangling tasks.

  • clock – Clock implementation for timing operations.

  • conversation_id – Optional custom ID for the conversation.

  • enable_rtvi – Whether to automatically add RTVI support to the pipeline. None, the default, adds it unless the pipeline is bridged: a bridged worker has no client of its own, and its RTVI would report every frame a second time as it crosses the bridge.

  • enable_tracing – Whether to enable tracing.

  • enable_turn_tracking – Whether to enable turn tracking.

  • exclude_frames – When bridged is set, extra frame types that should not cross the bus (lifecycle frames are always excluded).

  • idle_timeout_frames – A tuple with the frames that should trigger an idle timeout if not received within idle_timeout_secs. The default pairs the VAD-only UserSpeakingFrame with the turn and transcription frames a provider-driven pipeline reports instead.

  • idle_timeout_secs – Timeout (in seconds) to consider pipeline idle or None. If a pipeline is idle the pipeline worker will be cancelled automatically.

  • name – Optional worker name (used for worker-style addressing on the bus).

  • observers – List of observers for monitoring pipeline execution.

  • processor_unusable_policy – What to do when a processor reports an error that leaves it unable to do its job, such as a service whose API key was rejected. Defaults to ProcessorUnusablePolicy.CONTINUE, leaving the decision to on_pipeline_error handlers.

  • params – Configuration parameters for the pipeline.

  • rtvi_observer_params – The RTVI observer parameter to use if RTVI is enabled.

  • rtvi_processor – The RTVI processor to add if RTVI is enabled.

  • setup_timeout_secs – Timeout (in seconds) to wait for every processor to be set up. Processors connect while they are set up, so one that blocks connecting never lets the pipeline start, and the worker is torn down instead of waiting forever.

  • start_timeout_secs – Timeout (in seconds) to wait for the StartFrame to reach the end of the pipeline. A processor that blocks while handling it never lets the pipeline start, so the worker is torn down instead of waiting forever.

  • task_manager – Optional task manager for handling asyncio tasks.

  • tool_resources –

    Deprecated alias for app_resources.

    Deprecated since version 1.2.0: Use app_resources instead. tool_resources will be removed in 2.0.0.

property params: PipelineParams

Get the pipeline parameters for this worker.

Returns:

The pipeline parameters configuration.

property bridged: bool

Whether this pipeline is bridged onto the bus.

property app_resources: Any

Get the application-defined resources passed to this worker.

This is the same object passed to the constructor as app_resources. Tool handlers can also access it via FunctionCallParams.app_resources. The framework returns the original reference; mutations are visible to all callers.

Returns:

The application-defined resources, or None if none were passed.

property pipeline: BasePipeline

Get the full pipeline managed by this pipeline worker.

This will also include any internal processors added by the pipeline worker.

Returns:

The pipeline managed by the pipeline worker.

property turn_tracking_observer: TurnTrackingObserver | None

Get the turn tracking observer if enabled.

Returns:

The turn tracking observer instance or None if not enabled.

property turn_trace_observer: TurnTraceObserver | None

Get the turn trace observer if enabled.

Returns:

The turn trace observer instance or None if not enabled.

property rtvi: RTVIProcessor

Get the RTVI processor if RTVI is enabled.

Returns:

The RTVI processor added to the pipeline when RTVI is enabled.

property reached_upstream_types: tuple[type[Frame], ...]

Get the currently configured upstream frame type filters.

Returns:

Tuple of frame types that trigger the on_frame_reached_upstream event.

property reached_downstream_types: tuple[type[Frame], ...]

Get the currently configured downstream frame type filters.

Returns:

Tuple of frame types that trigger the on_frame_reached_downstream event.

add_observer(observer: BaseObserver)[source]

Add an observer to monitor pipeline execution.

Parameters:

observer – The observer to add to the pipeline monitoring.

async remove_observer(observer: BaseObserver)[source]

Remove an observer from pipeline monitoring.

Parameters:

observer – The observer to remove from pipeline monitoring.

set_reached_upstream_filter(types: tuple[type[Frame], ...])[source]

Set which frame types trigger the on_frame_reached_upstream event.

Parameters:

types – Tuple of frame types to monitor for upstream events.

set_reached_downstream_filter(types: tuple[type[Frame], ...])[source]

Set which frame types trigger the on_frame_reached_downstream event.

Parameters:

types – Tuple of frame types to monitor for downstream events.

add_reached_upstream_filter(types: tuple[type[Frame], ...])[source]

Add frame types to trigger the on_frame_reached_upstream event.

Parameters:

types – Tuple of frame types to add to upstream monitoring.

add_reached_downstream_filter(types: tuple[type[Frame], ...])[source]

Add frame types to trigger the on_frame_reached_downstream event.

Parameters:

types – Tuple of frame types to add to downstream monitoring.

has_finished() → bool[source]

Check if the pipeline worker has finished execution.

This indicates whether the worker has finished, meaning all processors have stopped.

Returns:

True if all processors have stopped and the worker is complete.

async stop_when_done()[source]

Schedule the pipeline to stop after processing all queued frames.

Sends an EndFrame to gracefully terminate the pipeline once all current processing is complete.

async end(*, reason: str | None = None) → None[source]

Request a graceful end of the session, draining the pipeline first.

Whatever this worker has already pushed reaches the end of the pipeline before the session goes away, rather than being cut off wherever it happened to be.

Parameters:

reason – Optional human-readable reason for ending.

async activate_worker(worker_name: str, *, args: WorkerActivationArgs | None = None, deactivate_self: bool = False) → None[source]

Activate a worker by name, draining this pipeline when handing over.

Deactivating this worker before its pipeline drains would let the target start producing while this worker’s output is still in flight, and both would arrive interleaved. A worker that stays active is handing nothing over, so it doesn’t wait: the first activation of a session would otherwise wait on the very worker it is about to wake.

Parameters:
  • worker_name – The name of the worker to activate.

  • args – Optional WorkerActivationArgs forwarded to the target worker’s on_activated.

  • deactivate_self – Whether to deactivate this worker before activating the target. Deactivating this worker drains its pipeline first; staying active does not. A worker that stays active and wants to drain anyway can call flush_pipeline() before this.

async cancel(*, reason: str | None = None)[source]

Request the running pipeline to cancel.

Parameters:

reason – Optional reason to indicate why the pipeline is being cancelled.

async run(params: WorkerParams)[source]

Start and manage the pipeline execution until completion or cancellation.

Parameters:

params – Configuration parameters for pipeline execution.

async queue_frame(frame: Frame, direction: FrameDirection = FrameDirection.DOWNSTREAM)[source]

Queue a single frame to be pushed through the pipeline.

Downstream frames are pushed from the beginning of the pipeline. Upstream frames are pushed from the end of the pipeline.

Parameters:
  • frame – The frame to be processed.

  • direction – The direction to push the frame. Defaults to downstream.

async queue_frames(frames: Iterable[Frame] | AsyncIterable[Frame], direction: FrameDirection = FrameDirection.DOWNSTREAM)[source]

Queue multiple frames to be pushed through the pipeline.

Downstream frames are pushed from the beginning of the pipeline. Upstream frames are pushed from the end of the pipeline.

Parameters:
  • frames – An iterable or async iterable of frames to be processed.

  • direction – The direction to push the frames. Defaults to downstream.

async flush_pipeline(timeout: float = 5.0) → bool[source]

Flush all in-flight frames from the pipeline and wait for it to drain.

Pushes a PipelineFlushFrame downstream; the sink bounces it back upstream, the source turns it around, and its event is set when it reaches the sink a second time. By then every frame queued ahead of it has been processed, along with anything a processor started by pushing upstream. The probe goes on the worker’s push queue, behind whatever is already waiting there, and bypasses any queue_frame override (e.g. tool-call deferral).

Parameters:

timeout – Seconds of no progress before giving up. Progress is a frame reaching this worker’s sink, or a report from the pipeline answering the probe when it crossed into another worker. A pipeline still working keeps the wait alive, however long it takes; one that is stuck, or that nobody will answer, gives up after this much quiet. On giving up a warning is logged and False is returned rather than blocking forever.

Returns:

True if the pipeline drained, False if it went quiet first.

track_flush_probe(frame: PipelineFlushFrame) → None[source]

Report progress on a probe from another worker while we hold it.

Called by whoever brings the probe into this pipeline, since only they know it is really coming in: every worker on the bus sees it, but most of them have nowhere to put it.

Parameters:

frame – The flush probe entering this pipeline.

async on_bus_message(message: BusMessage) → None[source]

Handle outbound bus messages: TTS playback and RTVI UI translation.

Runs the base lifecycle/job dispatch first. A BusTTSSpeakMessage targeted at this worker is queued as a TTSSpeakFrame (pipelines without a TTS service let it flow through). When this worker owns the RTVI processor, UI carriers produced by a UIWorker (BusUIDataMessage subclasses) are translated into RTVI frames by _handle_ui_bus_message; other workers skip the translation.

class pipecat.pipeline.worker.PipelineTask(**kwargs)[source]

Bases: PipelineWorker

Deprecated alias for PipelineWorker.

Deprecated since version 1.3.0: Use PipelineWorker instead. PipelineTask will be removed in 2.0.0.

class pipecat.pipeline.worker.PipelineTaskParams(**kwargs)[source]

Bases: WorkerParams

Deprecated alias for WorkerParams.

Deprecated since version 1.3.0: Use WorkerParams instead. Will be removed in 2.0.0.