Integracja z Pipecat

Użyj pipeline’u Pipecat jako mózgu LLM dla Speech Engine.

Ten przewodnik pokazuje, jak używać Pipecat jako pipeline’u LLM w serwerze brain dla Speech Engine. Speech Engine obsługuje pętlę głosową — zamianę mowy na tekst, zarządzanie turami i zamianę tekstu na mowę — a Pipecat generuje tekst przez modułowy pipeline procesorów (wywołania LLM, RAG, wywołania funkcji, guardraile, filtry treści).

Ten przewodnik dotyczy tylko Pythona, ponieważ Pipecat jest frameworkiem Pythona po stronie serwera. Nie ma odpowiednika dla Node dla procesorów pipeline’u; istnieje pakiet pipecat-client-js, ale jest to klient przeglądarkowy komunikujący się z serwerem Pipecat, a nie narzędzie do budowania pipeline’ów w TypeScript.

Architektura

SDK Speech Engine działa jako zewnętrzna warstwa — jego callback on_transcript uruchamia się za każdym razem, gdy użytkownik skończy mówić. W callbacku budujesz pipeline Pipecat, przekazujesz historię rozmowy jako LLMContextFrame i streamujesz tekst z pipeline’u z powrotem do Speech Engine. ElevenLabs zamienia tekst na mowę i odtwarza ją użytkownikowi.

User speaks (audio) on_transcript(history) LLMContextFrame(history) LLMTextFrame chunks send_response(async iterator) Agent speaks (audio) Browser ElevenLabs Brain Server (engine.serve) Pipecat Pipeline

Pipeline Pipecat działa tylko przez czas jednej tury. Gdy pojawi się nowa transkrypcja, poprzedni pipeline jest anulowany przed uruchomieniem kolejnego — w ten sposób obsługa przerwań w Speech Engine trafia do pipeline’u.

Kiedy użyć tego wzorca

Pipecat sprawdza się, gdy brain potrzebuje czegoś więcej niż pojedynczego wywołania LLM:

  • Modułowe procesory do generowania wspomaganego wyszukiwaniem, wywołań funkcji lub guardraili
  • Middleware oparte na ramkach, które może sprawdzać, przekształcać lub blokować ruch na każdym etapie
  • Fragmenty pipeline’u do ponownego użycia w wielu agentach

Jeśli twój brain to „transkrypcja na wejściu, wywołanie LLM na wyjściu”, prostszy będzie krótki przewodnik po Speech Engine. Sięgnij po Pipecat, gdy najważniejszy jest sam pipeline.

Wymagania

Zainstaluj zależności

pip install "pipecat-ai[openai]" "elevenlabs" "python-dotenv"

pipecat-ai[openai] instaluje usługę OpenAI LLM. Jeśli wolisz innego dostawcę, zamień rozszerzenie na odpowiednie (anthropic, google itd.).

Zbuduj brain Pipecat

Brain składa się z dwóch części: procesora TextSink, który zapisuje streamowany tekst do asyncio.Queue, oraz korutyny run_pipecat_brain, która buduje pipeline na jedną turę i zwraca fragmenty jako iterator asynchroniczny.

brain.py
import asyncio
import os
from typing import AsyncIterator
from dotenv import load_dotenv
from pipecat.frames.frames import (
Frame,
LLMContextFrame,
LLMFullResponseEndFrame,
LLMTextFrame,
)
from pipecat.pipeline.pipeline import Pipeline
from pipecat.pipeline.runner import PipelineRunner
from pipecat.pipeline.task import PipelineTask
from pipecat.processors.aggregators.llm_context import LLMContext
from pipecat.processors.frame_processor import FrameDirection, FrameProcessor
from pipecat.services.openai.llm import OpenAILLMService
load_dotenv()
SYSTEM_PROMPT = (
"You are a helpful voice assistant. Keep responses concise and conversational."
)
class TextSink(FrameProcessor):
"""Drain LLMTextFrame text into an asyncio.Queue."""
def __init__(self, queue: asyncio.Queue):
super().__init__()
self._queue = queue
async def process_frame(self, frame: Frame, direction: FrameDirection):
await super().process_frame(frame, direction)
if isinstance(frame, LLMTextFrame):
await self._queue.put(frame.text)
elif isinstance(frame, LLMFullResponseEndFrame):
await self._queue.put(None) # sentinel
await self.push_frame(frame, direction)
def build_messages(transcript: list[dict]) -> list[dict]:
messages = [{"role": "system", "content": SYSTEM_PROMPT}]
for turn in transcript:
role = "assistant" if turn["role"] == "agent" else turn["role"]
messages.append({"role": role, "content": turn["content"]})
return messages
async def run_pipecat_brain(transcript: list[dict]) -> AsyncIterator[str]:
"""Yield response text chunks from a one-turn Pipecat pipeline."""
llm = OpenAILLMService(
api_key=os.environ["OPENAI_API_KEY"],
model="gpt-4o-mini",
)
queue: asyncio.Queue[str | None] = asyncio.Queue()
sink = TextSink(queue)
task = PipelineTask(Pipeline([llm, sink]))
runner = PipelineRunner(handle_sigint=False)
async def drive():
context = LLMContext(build_messages(transcript))
await task.queue_frame(LLMContextFrame(context))
await task.stop_when_done()
run_task = asyncio.create_task(runner.run(task))
drive_task = asyncio.create_task(drive())
try:
while True:
chunk = await queue.get()
if chunk is None:
break
yield chunk
finally:
await task.cancel()
await asyncio.gather(run_task, drive_task, return_exceptions=True)

Pipeline zawiera tylko usługę LLM i sink — bez procesorów STT ani TTS, ponieważ obsługuje je Speech Engine. LLMContextFrame to dane wejściowe, a fragmenty LLMTextFrame to dane wyjściowe.

run_pipecat_brain jest generatorem asynchronicznym. Każdy zwrócony fragment trafia prosto do Speech Engine, więc agent zaczyna mówić, zanim cała odpowiedź będzie gotowa.

Podłącz go do serwera Speech Engine

Metoda send_response w SDK Speech Engine przyjmuje string lub dowolny asynchroniczny iterowalny obiekt stringów, więc możesz przekazać bezpośrednio run_pipecat_brain(transcript). Przed przekazaniem ich do brain zamień obiekty ConversationMessage z Speech Engine na zwykłe dicty.

server.py
import asyncio
import os
from dotenv import load_dotenv
from elevenlabs import AsyncElevenLabs
from brain import run_pipecat_brain
load_dotenv()
elevenlabs = AsyncElevenLabs(api_key=os.environ["ELEVENLABS_API_KEY"])
SPEECH_ENGINE_ID = os.environ["SPEECH_ENGINE_ID"]
async def on_transcript(transcript, session):
history = [{"role": m.role, "content": m.content} for m in transcript]
await session.send_response(run_pipecat_brain(history))
async def main():
engine = await elevenlabs.speech_engine.get(SPEECH_ENGINE_ID)
await engine.serve(
port=3001,
path="/ws",
debug=True,
on_transcript=on_transcript,
)
if __name__ == "__main__":
asyncio.run(main())

SDK Speech Engine anuluje zadanie poprzedniej tury, gdy pojawi się nowa transkrypcja. To anuluje generator asynchroniczny i bazowy PipelineTask przez blok try/finally w run_pipecat_brain.

Uruchom serwer

ngrok http 3001
python server.py

Połącz się ze Speech Engine z przeglądarki, używając tego samego endpointu tokena i kodu klienta co w krótkim przewodniku. Pipeline Pipecat działa po stronie serwera, a przeglądarka widzi zwykłą rozmowę ze Speech Engine.

Rozbuduj pipeline

Pipeline Pipecat obsługujący tylko tekst może zawierać dowolny procesor ramek, który działa na LLMTextFrame lub LLMContextFrame. Oto kilka częstych dodatków:

  • Guardraile: FrameProcessor umieszczony przed LLM, który sprawdza LLMContextFrame oraz zastępuje lub blokuje niebezpieczny kontekst.
  • Wywołania funkcji: zarejestruj narzędzia w OpenAILLMService, a Pipecat natywnie obsłuży ramki wywołań narzędzi. Końcowy tekst asystenta nadal trafia jako LLMTextFrame.
  • Wielostopniowe rozumowanie: połącz dwie instancje OpenAILLMService, z własnym procesorem między nimi, który przepisuje kontekst dla drugiego przebiegu.
  • Filtrowanie wyjścia: FrameProcessor umieszczony po LLM, który sprawdza każdy LLMTextFrame i odrzuca lub przepisuje niedozwolone treści, zanim dotrą do TextSink.

Kształt pipeline’u pozostaje taki sam — Pipeline([processor_a, llm, processor_b, sink]) — a run_pipecat_brain się nie zmienia.

Uwagi produkcyjne

  • Bezpieczne anulowanie: PipelineTask.cancel() może się zablokować, jeśli zostanie wywołane, zanim pipeline w pełni się uruchomi (pipecat-ai/pipecat#4276). Powyższy wzorzec try/finally jest bezpieczny, ponieważ cancel() uruchamia się dopiero po zakolejkowaniu co najmniej jednej ramki.
  • Prompt injection: wynik zamiany mowy na tekst to dane wejściowe użytkownika. Sprawdź lub znormalizuj transkrypcję przed przekazaniem jej do LLM, zwłaszcza jeśli procesor niższego poziomu używa tekstu w wywołaniach narzędzi lub zapytaniach do bazy danych.
  • Uwierzytelnianie serwera brain: ustaw wspólny sekret w Speech Engine i sprawdzaj go na serwerze brain, aby zapobiec nieautoryzowanym połączeniom z endpointem /ws:
    await elevenlabs.speech_engine.update(
    speech_engine_id=SPEECH_ENGINE_ID,
    speech_engine={"request_headers": {"x-api-key": os.environ["SHARED_SECRET"]}},
    )
  • Dostawca LLM: pipecat-ai[openai] zawiera OpenAILLMService. W przypadku Anthropic zainstaluj pipecat-ai[anthropic] i użyj AnthropicLLMService; reszta pipeline’u pozostaje bez zmian.

Kolejne kroki