service_metrics_observer

Observer reporting what each service spent and consumed.

Services report metrics as they finish a piece of work, and this observer turns each one into a record: what was measured, which processor and model reported it, and when. Nothing is summed, so a consumer groups the records by turn, session or model as it needs, and a session that ends abruptly still leaves behind everything that happened before it did.

class pipecat.observers.service_metrics_observer.ServiceLatencyKind(*values)[source]

Bases: StrEnum

Which measurement of a service’s own time a record carries.

TTFB = 'ttfb'
TTFA = 'ttfa'
TTFAT = 'ttfat'
class pipecat.observers.service_metrics_observer.ServiceUsageKind(*values)[source]

Bases: StrEnum

Which kind of service consumed something.

STT = 'stt'
LLM = 'llm'
TTS = 'tts'
class pipecat.observers.service_metrics_observer.ServiceLatencyRecord(*, kind: ServiceLatencyKind, processor: str, model: str | None = None, timestamp: float, seconds: float, ttfb_secs: float | None = None, leading_silence_secs: float | None = None, thinking_time_secs: float | None = None)[source]

Bases: BaseModel

One measurement of how long a service took.

Parameters:
  • kind – Which wait was measured.

  • processor – Name of the processor that reported it.

  • model – Model the processor was using, where it names one.

  • timestamp – Unix timestamp when the measurement was observed.

  • seconds – The measurement itself.

  • ttfb_secs – The time to first byte the measurement builds on, for the kinds that report one.

  • leading_silence_secs – Silence at the head of the first audio, for time to first audio.

  • thinking_time_secs – Time between a model’s first output and its first answer token, for time to first answer token.

class pipecat.observers.service_metrics_observer.ServiceUsageRecord(*, kind: ServiceUsageKind, processor: str, model: str | None = None, timestamp: float, audio_seconds: float | None = None, characters: int | None = None, prompt_tokens: int | None = None, completion_tokens: int | None = None, total_tokens: int | None = None, cache_read_input_tokens: int | None = None, cache_creation_input_tokens: int | None = None, reasoning_tokens: int | None = None, input_audio_tokens: int | None = None, output_audio_tokens: int | None = None, cache_read_input_audio_tokens: int | None = None)[source]

Bases: BaseModel

What one service consumed doing a piece of work.

A field is set only where the kind of service reports it, so an LLM record carries token counts and a text-to-speech record carries characters.

Parameters:
  • kind – Which kind of service reported.

  • processor – Name of the processor that reported it.

  • model – Model the processor was using, where it names one.

  • timestamp – Unix timestamp when the usage was observed.

  • audio_seconds – Audio transcribed, for speech-to-text.

  • characters – Characters synthesised, for text-to-speech.

  • prompt_tokens – Tokens in the prompt, for an LLM.

  • completion_tokens – Tokens generated, for an LLM.

  • total_tokens – Tokens in the prompt and the completion together.

  • cache_read_input_tokens – Prompt tokens served from cache.

  • cache_creation_input_tokens – Prompt tokens written to cache.

  • reasoning_tokens – Tokens spent reasoning before answering.

  • input_audio_tokens – Audio tokens in the prompt.

  • output_audio_tokens – Audio tokens generated.

  • cache_read_input_audio_tokens – Audio prompt tokens served from cache.

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

Bases: BaseObserver

Reports each metric a service publishes as its own record.

A record arrives per piece of work rather than per turn or per session: a turn that runs two inferences reports two, and a consumer that wants a total groups them itself. Summing here would lose the grain, and a total held in memory is lost with the process holding it.

What a service made someone wait for is here; what it did with its own time is not. Processing time, text aggregation and smart-turn predictions are all deliberately absent: aggregation already appears as a span in LatencyBreakdown, and the other two describe how work was done rather than what it cost the person waiting.

Events:
on_service_latency(observer, record): Emitted for each measurement of

a service’s own time, as a ServiceLatencyRecord.

on_service_usage(observer, record): Emitted for each report of what a

service consumed, as a ServiceUsageRecord.

Example:

observer = ServiceMetricsObserver()

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

Initialize the service metrics observer.

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

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

async on_push_frame(data: FramePushed)[source]

Report the metrics carried by a frame.

Metrics travel in one direction. A frame broadcast both ways would arrive as two frames, and would be reported twice.

Parameters:

data – Frame push event containing the frame and direction.