Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 11 additions & 4 deletions py/src/braintrust/integrations/pipecat/test_pipecat.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,11 +12,11 @@
from braintrust import SpanCustomizer, logger, set_span_customizers
from braintrust.integrations.pipecat import (
BraintrustPipecatObserver,
PipecatIntegration,
setup_pipecat,
wrap_pipeline_worker,
)
from braintrust.integrations.test_utils import verify_autoinstrument_script
from braintrust.integrations.versioning import detect_module_version, version_satisfies
from braintrust.logger import Attachment
from braintrust.test_helpers import init_test_logger

Expand Down Expand Up @@ -155,7 +155,6 @@ def _worker_runner_kwargs(**overrides):
return kwargs


@pytest.mark.vcr
@pytest.mark.asyncio
async def test_pipecat_observer_capture_audio_attachments_adds_tts_and_user_audio(memory_logger):
TTSStartedFrame = _import("pipecat.frames.frames.TTSStartedFrame")
Expand Down Expand Up @@ -297,16 +296,25 @@ async def test_setup_pipecat_traces_real_pipeline_frames(memory_logger):

@worker.event_handler("on_pipeline_started")
async def on_pipeline_started(_worker, _frame):
await worker.queue_frames([LLMContextFrame(context), EndFrame()])
await worker.queue_frames([LLMContextFrame(context), EndFrame(reason="pipeline complete")])

runner = WorkerRunner(**_worker_runner_kwargs())
await runner.add_workers(worker)
await asyncio.wait_for(runner.run(), timeout=20)

observer = next(o for o in getattr(worker, "_observer")._observers if isinstance(o, BraintrustPipecatObserver))
if version_satisfies(detect_module_version(importlib.import_module("pipecat"), ("pipecat",)), ">=1.12.0"):
assert observer.observe_every_push is False
assert not hasattr(observer, "_seen_frame_ids")
else:
assert observer._seen_frame_ids

logs = memory_logger.pop()
pipeline_span = _single_span(logs, "pipecat_pipeline")
assert _span_type(pipeline_span) == "task"
assert pipeline_span.get("metrics", {}).get("end") is not None
assert pipeline_span["metadata"]["terminal_frame"] == "EndFrame"
assert pipeline_span["metadata"]["reason"] == "pipeline complete"

llm_span = _single_span(logs, "pipecat_llm_response")
assert _span_type(llm_span) == "task"
Expand Down Expand Up @@ -336,7 +344,6 @@ def test_setup_and_wrap_pipeline_worker_are_idempotent():
PipelineWorker = _import("pipecat.pipeline.worker.PipelineWorker")
IdentityFilter = _import("pipecat.processors.filters.identity_filter.IdentityFilter")

assert PipecatIntegration.min_version == "1.3.0"
assert setup_pipecat(project_name="test-project-pipecat-py-tracing")
assert setup_pipecat(project_name="test-project-pipecat-py-tracing")
assert setup_pipecat(project_name="test-project-pipecat-py-tracing", capture_audio_attachments=True)
Expand Down
21 changes: 18 additions & 3 deletions py/src/braintrust/integrations/pipecat/tracing.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@
_pcm_to_wav,
_resolve_audio_attachment_options,
)
from braintrust.integrations.versioning import detect_module_version, version_satisfies
from braintrust.logger import NOOP_SPAN, Attachment, SpanTypeAttribute, current_span
from braintrust.logger import start_span as _bt_start_span

Expand All @@ -24,8 +25,10 @@ def start_span(*args, **kwargs):


try:
import pipecat
from pipecat.observers.base_observer import BaseObserver
except ImportError: # pragma: no cover - exercised when Pipecat is not installed.
_USES_NATIVE_FRAME_DEDUPLICATION = False

class BaseObserver: # type: ignore[no-redef]
"""Fallback base so this module is importable without Pipecat."""
Expand All @@ -35,6 +38,9 @@ def __init__(self, **_kwargs: Any) -> None:

async def cleanup(self) -> None:
pass
else:
_pipecat_version = detect_module_version(pipecat, ("pipecat",))
_USES_NATIVE_FRAME_DEDUPLICATION = version_satisfies(_pipecat_version, ">=1.12.0")


_TERMINAL_FRAME_TYPES = {"EndFrame", "StopFrame", "CancelFrame"}
Expand Down Expand Up @@ -74,6 +80,9 @@ def __init__(
trace_turns: bool = True,
**kwargs: Any,
) -> None:
self._uses_native_frame_deduplication = _USES_NATIVE_FRAME_DEDUPLICATION
if self._uses_native_frame_deduplication:
kwargs["observe_every_push"] = False
super().__init__(**kwargs)
(
self.capture_user_audio_attachments,
Expand All @@ -88,7 +97,8 @@ def __init__(
self._parent = _current_parent_export()
self._pipeline_span: Any | None = None
self._pipeline_parent: str | None = None
self._seen_frame_ids: set[int] = set()
if not self._uses_native_frame_deduplication:
self._seen_frame_ids: set[int] = set()
self._latest_llm_input: Any = None
self._latest_llm_metadata: dict[str, Any] = {}
self._llm_span: Any | None = None
Expand All @@ -111,7 +121,12 @@ async def on_pipeline_started(self) -> None:
self._ensure_pipeline_span()

async def on_process_frame(self, data: Any) -> None:
await self._handle_frame(getattr(data, "frame", None), processor=getattr(data, "processor", None))
frame = getattr(data, "frame", None)
processor = getattr(data, "processor", None)
is_terminal_at_sink = type(frame).__name__ in _TERMINAL_FRAME_TYPES and _is_pipeline_sink_processor(processor)
if self._uses_native_frame_deduplication and not is_terminal_at_sink:
return
Comment on lines 123 to +128

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Preserve terminal-frame delivery at the pipeline sink

With Pipecat 1.12+, observe_every_push=False delivers a frame only on its first push, whose source is not the pipeline sink. The terminal-frame guard at lines 139–140 therefore discards every EndFrame, StopFrame, and CancelFrame, while this early return also disables the per-processor callback that could observe the frame at the sink. Consequently _end_pipeline_span() never records terminal_frame or reason, and the span remains open until generic observer cleanup; keep a sink-visible path for terminal frames while deduplicating the other frame types.

Useful? React with 👍 / 👎.

await self._handle_frame(frame, processor=processor)

async def on_push_frame(self, data: Any) -> None:
await self._handle_frame(getattr(data, "frame", None), processor=getattr(data, "source", None))
Expand All @@ -128,7 +143,7 @@ async def _handle_frame(self, frame: Any, *, processor: Any = None) -> None:
return

frame_id = getattr(frame, "id", None)
if isinstance(frame_id, int):
if isinstance(frame_id, int) and not self._uses_native_frame_deduplication:
if frame_id in self._seen_frame_ids:
return
self._seen_frame_ids.add(frame_id)
Expand Down
Loading