Intégration de Pipecat

Utilisez un pipeline Pipecat comme cerveau LLM derrière Speech Engine.

Ce guide explique comment utiliser Pipecat comme pipeline LLM au sein d’un serveur cerveau Speech Engine. Speech Engine gère la boucle vocale, soit la conversion de la parole en texte, la gestion des tours de parole et la conversion du texte en parole, tandis que Pipecat génère le texte via un pipeline composable de processeurs (appels LLM, RAG, appels de fonctions, garde-fous, filtres de contenu).

Ce guide est uniquement disponible en Python, car Pipecat est un framework Python côté serveur. Il n’existe pas d’équivalent Node pour les processeurs de pipeline ; un package pipecat-client-js existe, mais il s’agit d’un client de navigateur qui communique avec un serveur Pipecat, et non d’un moyen de créer des pipelines en TypeScript.

Architecture

Le SDK Speech Engine constitue la couche externe : son callback on_transcript se déclenche chaque fois que l’utilisateur finit de parler. Dans ce callback, vous créez un pipeline Pipecat, alimentez l’historique de la conversation sous forme de LLMContextFrame, puis transmettez en continu la sortie textuelle du pipeline à Speech Engine. ElevenLabs convertit le texte en parole et le diffuse à l’utilisateur.

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

Le pipeline Pipecat ne s’exécute que pendant un tour de parole. Lorsqu’une nouvelle transcription arrive, le pipeline précédent est annulé avant l’exécution du suivant : c’est ainsi que la gestion des interruptions de Speech Engine se propage au pipeline.

Quand utiliser ce modèle

Pipecat excelle lorsque votre cerveau a besoin de plus qu’un simple appel LLM :

  • Des processeurs composables pour la génération augmentée par récupération, les appels de fonctions ou les garde-fous
  • Un middleware basé sur des frames, capable d’inspecter, de transformer ou de bloquer le trafic à chaque étape
  • Des fragments de pipeline réutilisables et partagés entre plusieurs agents

Si votre cerveau se limite à « transcription en entrée, appel LLM en sortie », le guide de démarrage Speech Engine est plus simple. Choisissez Pipecat lorsque le pipeline lui-même est l’élément central.

Prérequis

Installer les dépendances

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

pipecat-ai[openai] installe le service LLM OpenAI. Remplacez cette extension par celle d’un autre fournisseur (anthropic, google, etc.) selon vos préférences.

Créer le cerveau Pipecat

Le cerveau comprend deux éléments : un processeur TextSink qui place le texte généré en continu dans une asyncio.Queue, et une coroutine run_pipecat_brain qui crée un pipeline à un tour et génère des fragments sous forme d’itérateur asynchrone.

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)

Le pipeline contient uniquement le service LLM et le récepteur, sans processeur STT ni TTS, car Speech Engine s’en charge. LLMContextFrame constitue l’entrée ; les fragments LLMTextFrame constituent la sortie.

run_pipecat_brain est un générateur asynchrone. Chaque fragment généré est transmis directement à Speech Engine ; l’agent commence donc à parler avant que la réponse complète soit prête.

Le connecter au serveur Speech Engine

La méthode send_response du SDK Speech Engine accepte une chaîne ou tout itérable asynchrone de chaînes. Vous pouvez donc transmettre directement run_pipecat_brain(transcript). Convertissez les objets ConversationMessage fournis par Speech Engine en dictionnaires simples avant de les transmettre au cerveau.

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

Le SDK Speech Engine annule la tâche du tour précédent lorsqu’une nouvelle transcription arrive, ce qui annule le générateur asynchrone et la PipelineTask sous-jacente grâce au bloc try/finally de run_pipecat_brain.

Exécuter le serveur

ngrok http 3001
python server.py

Connectez-vous à Speech Engine depuis un navigateur à l’aide du même endpoint de jeton et code client que dans le guide de démarrage. Le pipeline Pipecat s’exécute côté serveur ; le navigateur voit une conversation Speech Engine classique.

Étendre le pipeline

Un pipeline Pipecat uniquement textuel peut inclure tout processeur de frames qui opère sur LLMTextFrame ou LLMContextFrame. Voici quelques ajouts courants :

  • Garde-fous : un FrameProcessor placé avant le LLM, qui inspecte LLMContextFrame et remplace ou bloque le contexte non sécurisé.
  • Appels de fonctions : enregistrez des outils sur OpenAILLMService, et Pipecat gère nativement les frames d’appel d’outils. Le texte final de l’assistant arrive toujours sous forme de LLMTextFrame.
  • Raisonnement en plusieurs étapes : enchaînez deux instances OpenAILLMService, avec un processeur personnalisé entre les deux qui réécrit le contexte pour le second passage.
  • Filtrage de sortie : un FrameProcessor placé après le LLM, qui inspecte chaque LLMTextFrame et supprime ou réécrit le contenu non autorisé avant qu’il n’atteigne TextSink.

La structure du pipeline reste identique : Pipeline([processor_a, llm, processor_b, sink]) ; run_pipecat_brain ne change pas.

Considérations pour la production

  • Sécurité de l’annulation : PipelineTask.cancel() peut provoquer un blocage s’il est appelé avant le démarrage complet du pipeline (pipecat-ai/pipecat#4276). Le modèle try/finally ci-dessus est sûr, car cancel() ne s’exécute qu’après qu’au moins une frame a été mise en file d’attente.
  • Injection de prompt : la sortie de la conversion de la parole en texte est une entrée utilisateur. Validez ou normalisez la transcription avant de l’envoyer au LLM, en particulier si un processeur en aval utilise le texte dans des appels d’outils ou des requêtes de base de données.
  • Authentification du serveur cerveau : définissez un secret partagé sur Speech Engine et vérifiez-le dans le serveur cerveau afin d’empêcher les connexions non autorisées à votre endpoint /ws :
    await elevenlabs.speech_engine.update(
    speech_engine_id=SPEECH_ENGINE_ID,
    speech_engine={"request_headers": {"x-api-key": os.environ["SHARED_SECRET"]}},
    )
  • Fournisseur LLM : pipecat-ai[openai] inclut OpenAILLMService. Pour Anthropic, installez pipecat-ai[anthropic] et utilisez AnthropicLLMService ; le reste du pipeline ne change pas.

Prochaines étapes