Pipecat-Integration

Verwenden Sie eine Pipecat-Pipeline als LLM-Brain hinter Speech Engine.

Dieser Leitfaden zeigt, wie Sie Pipecat als LLM-Pipeline in einem Speech-Engine-Brain-Server verwenden. Speech Engine übernimmt die Sprachschleife — Speech to Text, Sprecherwechsel und Text to Speech — während Pipecat die Textgenerierung über eine kombinierbare Pipeline von Prozessoren übernimmt (LLM-Aufrufe, RAG, Funktionsaufrufe, Guardrails, Inhaltsfilter).

Dieser Leitfaden gilt nur für Python, da Pipecat serverseitig ein Python-Framework ist. Es gibt kein Node-Äquivalent für die Pipeline-Prozessoren. Zwar existiert ein Paket namens pipecat-client-js, es ist jedoch ein Browser-Client, der mit einem Pipecat-Server kommuniziert, und keine Möglichkeit, Pipelines in TypeScript zu erstellen.

Architektur

Das Speech Engine SDK fungiert als äußere Schicht — sein on_transcript-Callback wird jedes Mal ausgelöst, wenn der Nutzer nicht mehr spricht. Innerhalb des Callbacks erstellen Sie eine Pipecat-Pipeline, übergeben den Gesprächsverlauf als LLMContextFrame und streamen die Textausgabe der Pipeline zurück an Speech Engine. ElevenLabs wandelt den Text in Sprache um und spielt ihn dem Nutzer ab.

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

Die Pipecat-Pipeline läuft nur für die Dauer eines Sprecherwechsels. Wenn ein neues Transkript eintrifft, wird die vorherige Pipeline abgebrochen, bevor die nächste startet — so wird die Unterbrechungsbehandlung von Speech Engine an die Pipeline weitergegeben.

Wann Sie dieses Muster verwenden sollten

Pipecat ist ideal, wenn Ihr Brain mehr als einen einzelnen LLM-Aufruf benötigt:

  • Kombinierbare Prozessoren für Retrieval-Augmented Generation, Funktionsaufrufe oder Guardrails
  • Frame-basierte Middleware, die Traffic bei jedem Schritt prüfen, transformieren oder blockieren kann
  • Wiederverwendbare Pipeline-Fragmente, die von mehreren Agents genutzt werden

Wenn Ihr Brain nach dem Muster „Transkript rein, LLM-Aufruf raus“ arbeitet, ist der Speech Engine Quickstart einfacher. Verwenden Sie Pipecat, wenn die Pipeline selbst der interessante Teil ist.

Voraussetzungen

  • Eine Speech Engine. Folgen Sie dem Speech Engine Quickstart, um eine zu erstellen.
  • Python 3.10+ (erforderlich für pipecat-ai).
  • Öffentlicher HTTPS-Tunnel für den Brain-Server (z. B. ngrok).

Abhängigkeiten installieren

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

pipecat-ai[openai] installiert den OpenAI-LLM-Service. Tauschen Sie das Extra bei Bedarf gegen einen anderen Anbieter aus (anthropic, google usw.).

Das Pipecat-Brain erstellen

Das Brain besteht aus zwei Teilen: einem TextSink-Prozessor, der gestreamten Text in eine asyncio.Queue schreibt, und einer run_pipecat_brain-Coroutine, die eine Pipeline für einen Sprecherwechsel erstellt und Chunks als asynchronen Iterator ausgibt.

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)

Die Pipeline enthält nur den LLM-Service und den Sink — keine STT- oder TTS-Prozessoren, da Speech Engine diese übernimmt. LLMContextFrame ist die Eingabe, LLMTextFrame-Chunks sind die Ausgabe.

run_pipecat_brain ist ein asynchroner Generator. Jeder ausgegebene Chunk geht direkt an Speech Engine, sodass der Agent zu sprechen beginnt, bevor die vollständige Antwort bereit ist.

Mit dem Speech-Engine-Server verbinden

send_response des Speech Engine SDK akzeptiert einen String oder jedes asynchrone Iterable von Strings. Sie können daher run_pipecat_brain(transcript) direkt übergeben. Wandeln Sie die von Speech Engine bereitgestellten ConversationMessage-Objekte in einfache Dicts um, bevor Sie sie an das Brain übergeben.

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())

Das Speech Engine SDK bricht die Aufgabe des vorherigen Sprecherwechsels ab, wenn ein neues Transkript eintrifft. Dadurch werden über den try/finally-Block in run_pipecat_brain auch der asynchrone Generator und die zugrunde liegende PipelineTask abgebrochen.

Server starten

ngrok http 3001
python server.py

Verbinden Sie sich über denselben Token-Endpunkt und Client-Code, die im Quickstart gezeigt werden, aus einem Browser mit der Speech Engine. Die Pipecat-Pipeline läuft serverseitig; im Browser erscheint eine normale Speech-Engine-Konversation.

Pipeline erweitern

Eine reine Text-Pipecat-Pipeline kann jeden Frame-Prozessor enthalten, der mit LLMTextFrame oder LLMContextFrame arbeitet. Einige häufige Ergänzungen:

  • Guardrails: Ein FrameProcessor vor dem LLM, der LLMContextFrame prüft und unsicheren Kontext ersetzt oder blockiert.
  • Funktionsaufrufe: Registrieren Sie Tools beim OpenAILLMService; Pipecat verarbeitet Tool-Call-Frames nativ. Der finale Assistant-Text kommt weiterhin als LLMTextFrame an.
  • Mehrstufiges Reasoning: Verketten Sie zwei OpenAILLMService-Instanzen mit einem benutzerdefinierten Prozessor dazwischen, der den Kontext für den zweiten Durchlauf umschreibt.
  • Ausgabefilterung: Ein FrameProcessor nach dem LLM, der jeden LLMTextFrame prüft und unzulässige Inhalte verwirft oder umschreibt, bevor sie TextSink erreichen.

Die Form der Pipeline bleibt gleich — Pipeline([processor_a, llm, processor_b, sink]) — und run_pipecat_brain ändert sich nicht.

Hinweise für den Produktionseinsatz

  • Sicherheit bei Abbrüchen: PipelineTask.cancel() kann zu einem Deadlock führen, wenn es aufgerufen wird, bevor die Pipeline vollständig gestartet ist (pipecat-ai/pipecat#4276). Das obige try/finally-Muster ist sicher, weil cancel() erst ausgeführt wird, nachdem mindestens ein Frame in die Queue gestellt wurde.
  • Prompt Injection: Die Speech-to-Text-Ausgabe ist Nutzereingabe. Validieren oder normalisieren Sie das Transkript, bevor Sie es an das LLM übergeben — insbesondere, wenn nachgelagerte Prozessoren den Text in Tool-Aufrufen oder Datenbankabfragen verwenden.
  • Authentifizierung des Brain-Servers: Legen Sie ein gemeinsames Secret für die Speech Engine fest und prüfen Sie es im Brain-Server, um unbefugte Verbindungen zu Ihrem /ws-Endpunkt zu verhindern:
    await elevenlabs.speech_engine.update(
    speech_engine_id=SPEECH_ENGINE_ID,
    speech_engine={"request_headers": {"x-api-key": os.environ["SHARED_SECRET"]}},
    )
  • LLM-Anbieter: pipecat-ai[openai] enthält OpenAILLMService. Installieren Sie für Anthropic pipecat-ai[anthropic] und verwenden Sie AnthropicLLMService; der Rest der Pipeline bleibt unverändert.

Nächste Schritte